S3に置いたログをTiDB Cloud Lakeへ継続的に取り込んでみた
こんにちは、ゲームソリューション部のsoraです。
今回は、アプリがS3に出し続けるログを、TiDB Cloud LakeのAWS – SQS(S3)データソースとIntegration Taskで継続的に取り込めるか試してみました。
先に結論
- S3のイベント通知をSQSに流しておけば、コンソールの画面操作だけで、S3に置かれたファイルがLakeに自動で取り込まれた
- 取り込みの実体は、イベントで届いたファイル1つごとに実行される
COPY INTOだった - 取り込み後にS3のファイルを消す設定や、同じファイルを二重に取り込まない設定は、
COPY INTOのPURGE/FORCEにそのまま対応していた - 列が欠けた行や壊れた行があると、既定ではファイルごと失敗して再試行が続く。
continueにすると、その行だけ捨てて残りが入った
構成
今回の構成は以下です。

- ECS Fargateのアプリは、リクエストごとにJSONのログを1行、標準出力に書きます
- FireLens(Fluent Bit)がそれを10秒ごとにまとめて、S3に1行1JSONのファイルとして置きます
- S3はファイルができるたびに、
ObjectCreatedのイベントをSQSへ送ります - TiDB Cloud LakeはIAMロールを引き受けてSQSのイベントを受け取り、書かれているファイルをS3から読んで取り込みます
Lakeのデータソースのうち、S3に置いたファイルを取り込めるのは以下の3つです。
| データソース | 認証 | 取り込めるファイル |
|---|---|---|
AWS – Credentials |
アクセスキーのみ | S3のファイル |
TiDB |
IAMロール / アクセスキー | Dumplingが出力したファイルだけ |
AWS – SQS(S3) |
IAMロール | SQSのイベントで届いたS3のファイル |
このうち、アプリのログのファイルをIAMロールで取り込めるのはAWS – SQS(S3)でした。
TiDBデータソースもIAMロールで使えますが、読めるのはDumplingの出力だけです。
この点は以下の記事で扱っています。
AWS側の準備
AWS側はTerraformで作りました。
ここでは取り込みに関わる部分だけ載せます。
FireLensの出力
アプリのコンテナのlogConfigurationで、Fluent BitのS3出力を指定します。
logConfiguration = {
logDriver = "awsfirelens"
options = {
Name = "s3"
region = "ap-northeast-1"
bucket = aws_s3_bucket.logs.id
upload_timeout = "10s"
total_file_size = "5M"
use_put_object = "On"
s3_key_format = "/app-logs/%Y/%m/%d/%H%M%S-$UUID.json"
json_date_key = "false"
}
}
upload_timeoutの既定は10分なので、動きを追いやすいように10秒にしています。
ログが1行も無い間は、ファイルは作られません。
S3に置かれるファイルの中身は、1行に1つのJSONです。
ログルーター側のenable-ecs-log-metadataで、ecs_clusterやcontainer_nameなどのECSのメタデータも同じ行に足されます。
{"timestamp":"2026-09-25T03:29:52.310Z","level":"INFO","service":"lake-ingest-app","request_id":"20b41581-...","method":"GET","path":"/api/items","status":200,"latency_ms":10.95,"user_id":90010,"user_agent":"curl/8.7.1","message":"request completed","container_id":"...","container_name":"app","source":"stdout","ecs_cluster":"lake-ingest","ecs_task_arn":"arn:aws:ecs:ap-northeast-1:<アカウントID>:task/lake-ingest/...","ecs_task_definition":"lake-ingest:3"}
S3のイベント通知とSQS
SQSの標準キューを作り、S3のObjectCreatedイベントを送ります。
LakeのAWS – SQS(S3)データソースは標準キューにしか対応していないので、FIFOキューは使えません。
必要な設定は、以下の公式ドキュメントに載っています。
resource "aws_s3_bucket_notification" "logs" {
bucket = aws_s3_bucket.logs.id
queue {
queue_arn = aws_sqs_queue.s3_events.arn
events = ["s3:ObjectCreated:*"]
filter_prefix = "app-logs/"
filter_suffix = ".json"
}
depends_on = [aws_sqs_queue_policy.s3_events]
}
キューのポリシーでは、このバケットからのsqs:SendMessageだけを許可しています。
なお、通知を設定すると、S3からs3:TestEventというテスト用のメッセージが1件届きます。
ファイルを指していないメッセージなので、今回はLakeのタスクを作る前に手で消しておきました。
Lakeが引き受けるIAMロール
信頼ポリシーには、Lakeのコンソールに表示される2つのPlatform role ARNとExternal IDを入れます。
表示される場所は、データソースの作成画面です。
権限はS3の読み取りとSQSのメッセージの受信・削除です。
取り込み後にS3のファイルを消す設定を試すので、s3:DeleteObjectも付けています。
{
"Statement": [
{
"Effect": "Allow",
"Action": ["s3:ListBucket", "s3:GetBucketLocation"],
"Resource": "arn:aws:s3:::lake-ingest-<アカウントID>"
},
{
"Effect": "Allow",
"Action": ["s3:GetObject", "s3:DeleteObject"],
"Resource": "arn:aws:s3:::lake-ingest-<アカウントID>/app-logs/*"
},
{
"Effect": "Allow",
"Action": [
"sqs:ReceiveMessage",
"sqs:DeleteMessage",
"sqs:GetQueueAttributes",
"sqs:ChangeMessageVisibility"
],
"Resource": "arn:aws:sqs:ap-northeast-1:<アカウントID>:lake-ingest-s3-events"
}
]
}
データソースの登録
Data > Data Sources > Createで、ServiceにAWS – SQS(S3)を選びます。

Configの枠に出ている2つのPlatform role ARNとExternal IDが、IAMロールの信頼ポリシーに入れる値です。
入力する値はこちらです。
| 項目 | 値 |
|---|---|
| Role ARN | 先ほど作ったIAMロールのARN |
| Queue URL | SQSのキューのURL |
| Bucket Filter | バケット名 |
| Prefix Filter | app-logs/ |
| Suffix Filter | .json |
SQSのメッセージには、作られたファイルのバケット名とキーがそのまま入っています。
Lakeはプレフィックス配下を見に行くのではなく、メッセージに書かれたファイルだけを読みます。
Prefix Filter / Suffix Filterは、そのキーが条件に合うかで取り込むファイルを絞る設定です。
今回はS3の通知と同じ値にしています。
Test Connectivityが通ったらOKで保存します。
Integration Taskの作成
Data > Data Integrations > Createで、作ったデータソースを選びます。

File TypeはCSV / PARQUET / NDJSONの3つから選びます。
ファイルの拡張子は.jsonですが、中身が1行1JSONなのでNDJSONを選びます。
NDJSONにすると、CSV用の区切り文字やヘッダーの項目は消えます。
Advanced Optionsの3つは、ここでは既定のままにしておきます。
| 項目 | 既定値 | 内容 |
|---|---|---|
| Error Handling | abort |
取り込めない行があったときに止めるか(abort / continue) |
| Clean Up Original Files | Off | 取り込んだあとにS3のファイルを消すか |
| Allow Duplicate Imports | Off | 取り込み済みのファイルをもう一度取り込むか |
それぞれの動きは後半で確かめます。
Nextで進むと、プレビューが出ます。

このときS3には3つのファイルがありましたが、プレビューに出たのは2つだけでした。
出なかったのは、イベント通知を設定する前に作られたファイルです。
プレビューはS3の一覧ではなく、SQSに届いているイベントのファイルを見ています。
次の画面でWarehouseと取り込み先のテーブルを指定します。
データベースは先にWorksheetで作っておきました。
CREATE DATABASE IF NOT EXISTS ingest;

列と型はプレビューのデータから推定されます。
timestampはVARCHARで推定されていたので、TIMESTAMPに変えます。
JSONの中では"2026-09-25T02:36:20.351Z"のような文字列なので、時刻の列だとは判定されないためです。
TIMESTAMPに変えておけば、取り込み時にそのまま変換されました。
Createで作成すると、タスクはStoppedの状態で作られます。
取り込みの確認
Runを押すと、SQSに届いていた2件のイベントがそれぞれ1回ずつ実行されました。
タスクを作る前に届いていたイベントも、キューに残っていれば取り込まれます。

取り込んだあと、SQSのメッセージは消えていました。
Clean Up Original FilesはOffなので、S3のファイルは残っています。
裏で何が実行されているかは、Monitoring > SQL Historyで見られます。
ただし、Userを自分にしたままだと出てきません。
取り込みはsystem:serviceaccount:<テナントID>というユーザーで実行されていて、User Agentはlake-sqs-s3-consumerでした。

1つを開くと、中身はCOPY INTOでした。

COPY INTO `ingest`.`access_logs` (`method`, `latency_ms`, ..., `message`)
FROM 's3://lake-ingest-<アカウントID>/app-logs/2026/09/25/033258-efa5uDeg.json'
CONNECTION = (external_id = '***', role_arn = '***')
FILE_FORMAT = (type = NDJSON)
PURGE = true
FORCE = false
DISABLE_VARIANT_CHECK = false
ON_ERROR = abort
RETURN_FAILED_ONLY = false
FROMにはイベントで届いたファイルが1つだけ書かれていて、Stageは使わずにCONNECTIONで認証を直接渡しています。
画面の設定は、COPY INTOのオプションにそのまま対応していました。
| 画面の設定 | COPY INTOのオプション |
|---|---|
| Clean Up Original Files | PURGE |
| Allow Duplicate Imports | FORCE |
| Error Handling | ON_ERROR |
Warehouseは、タスクに指定したsora-blog-testで実行されていました。
ファイルが届くたびにWarehouseが起動するので、ログが流れ続ける環境ではWarehouseも動き続けることになります。
この状態でアプリを1回叩くと、S3にファイルができてすぐにCOPY INTOが実行され、行が増えました。
待ち時間はほぼFluent Bitがファイルをまとめる時間(upload_timeout)で、Lake側の待ちはほとんどありませんでした。

なお、timestampはUTCの値のまま保存されますが、Worksheetでは日本時間に直して表示されました。
S3のファイルの時刻と見比べるときは、9時間ずれて見えます。
タスクをStopしている間にアプリを叩いた分も、Startしたときに取り込まれました。
イベントはSQSに残っているので、SQSのメッセージの保持期間内なら、止めていても取りこぼしません。
Clean Up Original Filesの確認
タスクの編集画面で、Clean Up Original FilesをOnに変えます。
作成したあとでも変えられました。

この状態でアプリを叩くと、取り込まれたあとにS3のファイルが消えました。
COPY INTOもPURGE = trueで実行されています。
Offのときに取り込んだファイルは、S3に残ったままでした。
Onにしたあとに取り込んだファイルだけが消える動きです。
重複取り込みの確認
S3のイベント通知もSQSの標準キューも、同じメッセージが2回届くことがあります。
そこで、取り込み済みのファイルを同じキー・同じ中身でもう一度アップロードして、同じファイルのイベントをもう一度送ってみました。
aws s3 cp s3://lake-ingest-<アカウントID>/app-logs/2026/09/25/032952-L28IvVIt.json ./032952.json
aws s3 cp ./032952.json s3://lake-ingest-<アカウントID>/app-logs/2026/09/25/032952-L28IvVIt.json
COPY INTOは実行されましたが、Scan Rowsは0で、行は増えませんでした。

Allow Duplicate ImportsがOff(FORCE = false)なので、取り込み済みのファイルとしてスキップされています。
同じイベントが重複して届いても、二重には入りません。
このときClean UpはOnにしていましたが、スキップしたファイルはS3から消えませんでした。
PURGEが消すのは、その実行で取り込んだファイルだけです。
Clean UpでS3を空に保つつもりでも、スキップされたファイルは残り続けます。
取り込めない行があるときの確認
最後に、取り込めない行が混ざったファイルをS3に直接置いてみます。
Error Handlingがabortのとき
まず、テーブルの列の一部が無い行を置きました。
アプリのログから、ECSのメタデータの列を抜いた行です。

BadBytes. Code: 1046, Text = Missing value for column 4 (container_name String NULL). current FILE_FORMAT option: MISSING_FIELD_AS=ERROR
at file 'app-logs/2026/09/25/broken-test-01.json', line 0.
NDJSONの取り込みは、既定でMISSING_FIELD_AS=ERRORになっています。
テーブルの列が1つでも欠けた行があると、ファイルごと失敗します。
アプリの起動ログのように、リクエストのログとは項目が違う行が混ざる場合は、それだけで止まることになります。
失敗したメッセージはSQSから消されず、何度も再試行されていました。
Run HistoryのLast Runには、同じ実行のエラーが再試行のたびに更新されて残ります。
Error Handlingがcontinueのとき
ファイルを「全列そろった正しい行」と「JSONではない行」の2行に置き換えて、Error Handlingをcontinueに変えました。

ON_ERROR = continueでCOPY INTOが実行され、正しい行だけが入りました。
JSONではない行は捨てられています。
再試行を待っていたメッセージにも、変えた設定がそのまま使われました。
止めて気づきたいならabort、多少欠けても取り込み続けたいならcontinue、という選び方になりそうです。
補足:SQLで組む場合
SQSを使わずに、SQLのTaskで定期的にCOPY INTOして取り込む方法もあります。
公式ガイドでは、VectorでS3に置いたログをこの形で取り込んでいます。
AWS側はS3だけで済みますが、取り込むのはファイルができたときではなく、Taskに決めた間隔ごとになります。
最後に
今回は、アプリがS3に出し続けるログを、TiDB Cloud LakeのAWS – SQS(S3)データソースとIntegration Taskで継続的に取り込めるかを試してみました。
S3のイベント通知をSQSに流しておけば、画面操作だけでS3に置いたログを取り込み続けられました。
この記事がどなたかの参考になれば幸いです。




