S3に置いたログをTiDB Cloud Lakeへ継続的に取り込んでみた

S3に置いたログをTiDB Cloud Lakeへ継続的に取り込んでみた

S3に出力されるアプリログをTiDB Cloud Lakeで継続的に取り込むため、AWS – SQS(S3)データソースとIntegration Taskを試してみました。設定方法から実際の動作、エラーハンドリングまで、実装の流れを紹介します。
2026.09.25

こんにちは、ゲームソリューション部のsoraです。
今回は、アプリがS3に出し続けるログを、TiDB Cloud LakeのAWS – SQS(S3)データソースとIntegration Taskで継続的に取り込めるか試してみました。

先に結論

  • S3のイベント通知をSQSに流しておけば、コンソールの画面操作だけで、S3に置かれたファイルがLakeに自動で取り込まれた
  • 取り込みの実体は、イベントで届いたファイル1つごとに実行されるCOPY INTOだった
  • 取り込み後にS3のファイルを消す設定や、同じファイルを二重に取り込まない設定は、COPY INTOのPURGE / FORCEにそのまま対応していた
  • 列が欠けた行や壊れた行があると、既定ではファイルごと失敗して再試行が続く。
    continueにすると、その行だけ捨てて残りが入った

構成

今回の構成は以下です。

00-sr-architecture

  • 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の出力だけです。
この点は以下の記事で扱っています。

https://dev.classmethod.jp/articles/tidb-cloud-lake-migrate-from-tidb/

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キューは使えません。
必要な設定は、以下の公式ドキュメントに載っています。

https://docs.pingcap.com/tidbcloudlake/amazon-sqs-s3-iam-role/

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)を選びます。

01-sr-lake-sqs-datasource-form

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で、作ったデータソースを選びます。

02-sr-lake-integration-basic-info

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で進むと、プレビューが出ます。

03-sr-lake-integration-preview

このときS3には3つのファイルがありましたが、プレビューに出たのは2つだけでした。
出なかったのは、イベント通知を設定する前に作られたファイルです。
プレビューはS3の一覧ではなく、SQSに届いているイベントのファイルを見ています。

次の画面でWarehouseと取り込み先のテーブルを指定します。
データベースは先にWorksheetで作っておきました。

CREATE DATABASE IF NOT EXISTS ingest;

04-sr-lake-integration-target-timestamp-varchar

列と型はプレビューのデータから推定されます。
timestampはVARCHARで推定されていたので、TIMESTAMPに変えます。
JSONの中では"2026-09-25T02:36:20.351Z"のような文字列なので、時刻の列だとは判定されないためです。
TIMESTAMPに変えておけば、取り込み時にそのまま変換されました。

Createで作成すると、タスクはStoppedの状態で作られます。

取り込みの確認

Runを押すと、SQSに届いていた2件のイベントがそれぞれ1回ずつ実行されました。
タスクを作る前に届いていたイベントも、キューに残っていれば取り込まれます。

05-sr-lake-integration-run-history

取り込んだあと、SQSのメッセージは消えていました。
Clean Up Original FilesはOffなので、S3のファイルは残っています。

裏で何が実行されているかは、Monitoring > SQL Historyで見られます。
ただし、Userを自分にしたままだと出てきません。
取り込みはsystem:serviceaccount:<テナントID>というユーザーで実行されていて、User Agentはlake-sqs-s3-consumerでした。

06-sr-lake-sql-history-serviceaccount

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

07-sr-lake-copy-into-sql

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側の待ちはほとんどありませんでした。

08-sr-lake-databases-3rows

なお、timestampはUTCの値のまま保存されますが、Worksheetでは日本時間に直して表示されました。
S3のファイルの時刻と見比べるときは、9時間ずれて見えます。

タスクをStopしている間にアプリを叩いた分も、Startしたときに取り込まれました。
イベントはSQSに残っているので、SQSのメッセージの保持期間内なら、止めていても取りこぼしません。

Clean Up Original Filesの確認

タスクの編集画面で、Clean Up Original FilesをOnに変えます。
作成したあとでも変えられました。

09-sr-lake-integration-edit-cleanup-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で、行は増えませんでした。

10-sr-lake-copy-into-reupload-scan-0

Allow Duplicate ImportsがOff(FORCE = false)なので、取り込み済みのファイルとしてスキップされています。
同じイベントが重複して届いても、二重には入りません。

このときClean UpはOnにしていましたが、スキップしたファイルはS3から消えませんでした。
PURGEが消すのは、その実行で取り込んだファイルだけです。
Clean UpでS3を空に保つつもりでも、スキップされたファイルは残り続けます。

取り込めない行があるときの確認

最後に、取り込めない行が混ざったファイルをS3に直接置いてみます。

Error Handlingがabortのとき

まず、テーブルの列の一部が無い行を置きました。
アプリのログから、ECSのメタデータの列を抜いた行です。

11-sr-lake-copy-into-error-missing-field

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に変えました。

12-sr-lake-copy-into-on-error-continue

ON_ERROR = continueでCOPY INTOが実行され、正しい行だけが入りました。
JSONではない行は捨てられています。
再試行を待っていたメッセージにも、変えた設定がそのまま使われました。

止めて気づきたいならabort、多少欠けても取り込み続けたいならcontinue、という選び方になりそうです。

補足:SQLで組む場合

SQSを使わずに、SQLのTaskで定期的にCOPY INTOして取り込む方法もあります。
公式ガイドでは、VectorでS3に置いたログをこの形で取り込んでいます。

https://docs.pingcap.com/tidbcloudlake/ingest-json-logs-with-vector-cloud/

AWS側はS3だけで済みますが、取り込むのはファイルができたときではなく、Taskに決めた間隔ごとになります。

最後に

今回は、アプリがS3に出し続けるログを、TiDB Cloud LakeのAWS – SQS(S3)データソースとIntegration Taskで継続的に取り込めるかを試してみました。
S3のイベント通知をSQSに流しておけば、画面操作だけでS3に置いたログを取り込み続けられました。
この記事がどなたかの参考になれば幸いです。


TiDB Cloudの導入・サポートはクラスメソッドにお任せください

クラスメソッドでは、TiDB Cloudの導入から運用支援まで、豊富なノウハウでお客様をサポートしています。パフォーマンスの最適化やスケーラビリティに課題を抱えている方は、ぜひご相談ください。
詳細な導入事例やサービス内容について知りたい方は、こちらからご確認いただけます。

TiDB Cloudのサポート詳細を見る

この記事をシェアする

関連記事