Upgrade to Pro — share decks privately, control downloads, hide ads and more …

OpenSearch Data Prepperにおける「Iceberg からの CDC」開発を...

Sponsored · Your Podcast. Everywhere. Effortlessly. Share. Educate. Inspire. Entertain. You do you. We'll handle the rest.
Avatar for bering bering
August 06, 2026
59

OpenSearch Data Prepperにおける「Iceberg からの CDC」開発を通じた学び

Iceberg Meetup #7での登壇資料です
https://iceberg.connpass.com/event/399857/

Avatar for bering

bering

August 06, 2026

Transcript

  1. ⾃⼰紹介 ⽦⽥ 宗太郎 (ひきた そうたろう) ソリューションアーキテクト OpenSearch と Iceberg を軸に、

    ⼤規模データ活⽤に取り組むお客様を⽀援しています 趣味: ビッグデータに関わるオープンソース、 特に Apache Iceberg と OpenSearch の 調査、開発、コミュニティへの貢献 最近は、Lance が気になっている
  2. 本⽇の主題 • OpenSearch Data Prepper という ETL ツールに、 「Iceberg からの

    CDC を実現する機能を作ってみた経験の共有 • 特に、Iceberg ライブラリの IncrementalChangelogScan を深掘り - Iceberg の差分認識 (CDC) を実装したい⼈々の参考になるはず - Iceberg を深く理解して活⽤したいユーザーの参考にもなるかも オリジナルのソフトウェア Iceberg のライブラリ 3
  3. • OpenSearch は、Linux Foundation が開発 / 管理する、オープンソース の検索エンジン / プラットフォーム

    • 全⽂検索、ベクトル検索など、 幅広い機能を備える • Web サービスの検索バックエンドや ログ分析 / オブザーバビリティ、 AI ユースケースなどに⽤いられる • データ処理パイプライン (Data Prepper)、分析ダッシュボード (OpenSearch UI) など周辺系も充実
  4. データレイク / レイクハウスと OpenSearch の関係性 • データソースと OpenSearch の間のステージングレイヤとして活⽤ •

    ⽣データを前処理してから投⼊することも データソース データレイク 前処理 OpenSearch インデクシング処理
  5. 余談: Hadoop は検索エンジンと縁の深いプロジェクト 2002〜︓nutch は lucene のサブプロジェクトで、スケーラブルな Web クローラ 2004〜︓インターネット上のデータをスケーラブルに前処理するための仕組みと

    して、Google の 論⽂に着想を得て nutch に GFS と MapReduce が組み込まれる 2006〜︓Apache Hadoop として独⽴ https://en.wikipedia.org/wiki/Apache_Hadoop#History
  6. ストレージ上のデータを OpenSearch へ投⼊する⽅法 • ⼩規模データは Python 等のクライアントライブラリで投⼊すれば OK • ⼤規模データ

    / ⾼度な要件に対応する上では、Spark / Flink なども検討 • 継続的なデータ投⼊パイプラインを運⽤する上では、 OpenSearch Project 製のデータ収集ツールである Data Prepper も便利 Data Source Indexing OpenSearch
  7. Iceberg をデータソースとする場合の課題 ② 変更差分 (Insert, Update, Delete) を継続的に抽出して OpenSearch に投⼊する処理の実装難易度が⾼い

    - 検索システムでは、ソースの変更を継続的に反映したくなることがよくある - Insert のみであれば、スナップショット間の差分を抽出するだけで済むが、 障害ケースを想定した対応(チェックポイントなど)は⼀筋縄では⾏かない - Update, Delete まで扱おうとすると、差分認識⾃体もさらに複雑に… データソース 継続的な更新 Iceberg の更新差分を 認識して反映
  8. OpenSearch Data Prepper • • • • • Java 製、ストリーミングベースの

    ETL ソフト (OpenSearch と独⽴して稼働) データソースの変更を継続的に取り込み、主に OpenSearch へ同期する フィルタリング・変換・加⼯・集約もサポート 複数ノードでの分散処理によって、⼤規模なデータにもスケール可能 E2E Acknowledgement 機構により障害ケースでも整合性を確保できる データソース ターゲット Kafka Kinesis OpenTelemetry S3 File OpenSearch etc OpenSearch Prometheus S3 File Source Processor Sink Buffer etc https://docs.opensearch.org/latest/data-prepper/
  9. Data Prepper Iceberg CDC Source • Iceberg テーブルの変更を継続的に検出してターゲットに同期 - UPDATE,

    DELETE にも対応 (Copy on Write のみ) • 特定スナップショットのワンショット同期も可能 • スキーマ進化に対応 • Exactly Once でのデータ連携を保証 (OpenSearchがターゲットの場合) • Iceberg コアライブラリに実装がある各カタログ、ストレージの Iceberg に対応
  10. シンプルな YAML 定義でセットアップ pipeline-name: source: # ソース定義 iceberg: # Iceberg

    ソースを使⽤ tables: - table_name: "my_database.my_table" catalog: type: rest uri: "https://example.com/api/catalog" warehouse: "my_warehouse" polling_interval: "PT30S" identifier_columns: ["region", "category"] sink: # ターゲット定義 - opensearch: hosts: ["https://opensearch:9200"] document_id: "document_id" action: "${getMetadata(¥"bulk_action¥")}"
  11. IncrementalChangelogScan • 2 つのスナップショット間の差分を ChangelogScanTask のリストとして返す • ChangelogScanTask のリストとは、差分の反映に必要なファイルのリスト -

    operation: INSERT / DELETE (ファイルの追加 / 論理削除) - changeOrdinal︓スキャン範囲の中で、何番⽬のスナップショットに属するか - commitSnapshotId︓変更が⾏われたスナップショットの id • CoW のみサポート。MoR のサポートは開発中 (iceberg#14264) IncrementalChangelogScan scan = table.newIncrementalChangelogScan() .fromSnapshotInclusive(<始点となるスナップショットid>) .toSnapshot(<終点となるスナップショットid>); https://iceberg.apache.org/javadoc/latest/org/apache/iceberg/IncrementalChangelogScan.html
  12. IncrementalChangelogScan で INSERT 差分を把握する ① テーブルへの INSERT を⾏うコミットが 3回 発⽣

    S1: INSERT (1, Laptop, 1200), (2, Mouse, 25), (3, Keyboard, 75) → file a 追加 S2: INSERT (4, Monitor, 300), (5, Webcam, 90) → file b 追加 S3: INSERT (6, Cable, 15), (7, Charger, 40) → file c 追加 INSERT (8, Dock, 150), (9, Hub, 60) → file d 追加 ② S1 – S3 を対象に IncrementalChangelogScan 実⾏ [AddedRowsScanTask] operation=INSERT file=<file a> changeOrdinal=0 [AddedRowsScanTask] operation=INSERT file=<file b> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file c> changeOrdinal=2 [AddedRowsScanTask] operation=INSERT file=<file d> changeOrdinal=2 AddedRowsScanTask が⽰す各ファイルの各⾏を取り込めば OK
  13. Data Prepper Iceberg CDC Sourceの実装 (INSERT) ①定期ポーリングで最新スナップショットをチェック -> 更新がある場合②へ Leader

    Scheduler Node Changelog Worker Nodes ② IncrementalChangelogScan で読み 取りタスク(ファイルパス)の⼀覧を取得 ④ 各ワーカーでタスクを 分担してファイルを読み、 後続処理 (Buffer)へ流す ③ タスク⼀覧を Coordination Store に 書き込み ⑤ 全タスクの完了が確認でき たら、処理済スナップショッ ト id を記録 → ①へ戻る Coordination Store (以降 DB と表記) * 現時点では Coordination Store は Dynamo DB のみサポート
  14. IncrementalChangelogScan で DELETE 差分を把握する ① テーブルで DELETE が発⽣ S1: INSERT

    (1,Laptop,1200), (2,Mouse,25), (3,Keyboard,75) → file a 追加 S2: レコード (2,Mouse,25) を Copy on Write で DELETE → file a 論理削除 → (1,Laptop,1200), (3,Keyboard,75) を含む file b 追加 ② S1 – S2 を対象に IncrementalChangelogScan 実⾏ [AddedRowsScanTask] operation=INSERT file=<file a> changeOrdinal=0 [DeletedDataFileScanTask] operation=DELETE file=<file a> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file b> changeOrdinal=1 「どの⾏が削除されたか」を把握するには、file a と file b を突合して、 消えた⾏ / 消えていない⾏を特定する必要がある (Iceberg の世界ではこれを「carryover」と呼ぶ )
  15. IncrementalChangelogScan で UPDATE 差分を把握する ① テーブルで UPDATE が発⽣ S1: INSERT

    (1,Laptop,1200), (2,Mouse,25), (3,Keyboard,75) → file a 追加 S2: レコード (2,Mouse,25) を Copy on Write で UPDATE (price 25→30) → file a 論理削除 → (1,Laptop,1200), (2,Mouse,30), (3,Keyboard,75) を含む file b 追加 ② S1 – S2 を対象に IncrementalChangelogScan 実⾏ [AddedRowsScanTask] operation=INSERT file=<file a> changeOrdinal=0 [DeletedDataFileScanTask] operation=DELETE file=<file a> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file b> changeOrdinal=1 「どの⾏がどう変わったか」を把握するには、file a と file b を突合して、 変わった⾏ / 変わっていない⾏を特定する必要がある = Primary Key となるカラムの指定が必須
  16. 単⼀コミットで複数のファイルが更新された場合 ① テーブルで DELETE が発⽣ S1: INSERT (1,Laptop,1200), (2,Mouse,25), (3,Keyboard,75)

    → data file a 追加 INSERT (4, Monitor, 300), (5, Webcam, 90) → data file b 追加 S2: レコード (2,Mouse,25) と (4, Monitor, 300) を Copy on Write で DELETE → data file a 論理削除 → (1,Laptop,1200), (3,Keyboard,75) を含む data file c 追加 → data file b 論理削除 → (5, Webcam, 90) を含む data file d 追加 ② S1 – S2 を対象に IncrementalChangelogScan 実⾏ [AddedRowsScanTask] operation=INSERT file=<file a> changeOrdinal=0 [DeletedDataFileScanTask] operation=DELETE file=<file a> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file b> changeOrdinal=1 [DeletedDataFileScanTask] operation=DELETE file=<file c> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file d> changeOrdinal=1
  17. 単⼀コミットで複数のファイルが更新された場合 ① テーブルで DELETE が発⽣ IncrementalChangelogScan が返す changeOrdinal=1 S1: INSERT

    (1,Laptop,1200), (2,Mouse,25), (3,Keyboard,75)の各タスクは順不同 → data file a 追加 INSERT (4, Monitor, 300), (5, Webcam, 90) → data file b 追加 (List 順に反映しても整合性は保証されない︕) S2: レコード (2,Mouse,25) と (4, Monitor, を Copy on Write で DELETE → operation=DELETE file=<file a> と300) operation=DELETE file=<file a> → data file a 削除 file=<file c> と operation=INSERT file=<file d> operation=DELETE → (1,Laptop,1200), (3,Keyboard,75) を含む data file c 追加 が対応関係にあることは、ファイルの中⾝を⾒ないと分からない → data file b 削除 → (5, Webcam, 90) を含む data file d 追加 ② S1 – S2 を対象に IncrementalChangelogScan 実⾏ [AddedRowsScanTask] operation=INSERT file=<file a> changeOrdinal=0 [DeletedDataFileScanTask] operation=DELETE file=<file a> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file b> changeOrdinal=1 [DeletedDataFileScanTask] operation=DELETE file=<file c> changeOrdinal=1 [AddedRowsScanTask] operation=INSERT file=<file d> changeOrdinal=1
  18. どのように DELTE / UPDATE を反映するか id (PK), item, price の

    3 カラムで構成されるテーブルの場合… ⼊⼒ スナップショットで 発⽣した DELETE / INSERT 全⾏を取得 Step1 全カラムが⼀致する DELETE+INSERT を相殺 Step2 PK を元にペアとな るDELETE + INSERT を UPDATE として認識 Step3 残った DELETE はDELETEとして認識 DELETE 1, Laptop,1200 DELETE 2, Mouse,25 DELETE 3, Keyboard,75 DELETE 4, Monitor,75 INSERT 1, Laptop,1200 INSERT 2 , Mouse,30 INSERT 3 , Keyboard,75 DELETE 1, Laptop,1200 INSERT 1 , Laptop,1200 DELETE 3, Keyboard,75 INSERT 3 , Keyboard,75 DELETE 2, Mouse,25 DELETE 4, Monitor,75 INSERT 2 , Mouse,30 DELETE 2, Mouse,25 -> INSERT 2, Mouse,30 DELETE 4, Monitor,75 DELETE 4, Monitor,75 Step 1, Step 2 共に、同じ PK の DELETE と INSERT を突き合わせる処理 = PK ごとにレコードをグルーピングする必要がある 複数ノードでどう分散処理するか︖が課題に
  19. Data Prepper Iceberg CDC Sourceの実装 (UPDATE, DELETE) Leader Scheduler Node

    ①定期ポーリングで最新スナップショット をチェック -> 更新がある場合②へ ② IncrementalChangelogScanで タスク⼀覧を取得 -> タスクが DELETE を含む場合 SHUFFLE 処理へ (含まない場合は INSERT のみのフローへ)
  20. Leader Schedular Node ① IncrementalChangelogScan の 結果を WRITE タスクとして DB

    に登録 Node 1 ② WRITE タスクを DB から取得 し、ファイルを Read DELETE file a 1, Laptop 2, Mouse-1 3, Keyboard 4, Monitor Node 2 INSERT file b 1, Laptop 2, Mouse-2 3, Keyboard
  21. Leader Schedular Node ① IncrementalChangelogScan の 結果を WRITE タスクとして DB

    に登録 Node 1 ② WRITE タスクを DB から取得 し、ファイルを Read DELETE file a 1, Laptop 2, Mouse-1 3, Keyboard 4, Monitor ④ ③の完了後、DB からグループ情報を 抽出、⼩さいグループを束ねて READ タスクとして DB に登録 ③ PK のハッシュでグループに 分けてローカルディスクに Write → メタデータを DB に登録 グループ A 1 2 3 DELETE グループ B 4 DELETE Node 2 INSERT file b 1, Laptop 2, Mouse-2 3, Keyboard グループ A 1 2 3 INSERT グループ B (該当なし)
  22. Leader Schedular Node ① IncrementalChangelogScan の 結果を WRITE タスクとして DB

    に登録 Node 1 ② WRITE タスクを DB から取得 し、ファイルを Read DELETE file a 1, Laptop 2, Mouse-1 3, Keyboard 4, Monitor ④ ③の完了後、DB からグループ情報を 抽出、⼩さいグループを束ねて READ タスクとして DB に登録 ③ PK のハッシュでグループに 分けてローカルディスクに Write → メタデータを DB に登録 グループ A 1 2 3 DELETE グループ B 4 DELETE Node 2 ⑥ ⑤の完了後、処理済 スナップショット id を記録 → 最初へ戻る ⑤ READ タスクを DB から取得し、 担当グループに対応するデータを Pull → carryover / UPDATE を認識 グループ A 担当 1: D+I 2: D+I 3: D+I id 1, 3 は全カラム同じ値なので相殺 id 2 は Mouse-1 がMouse-2 になっ ているので UPDATE として後続反映 グループ B 担当 INSERT file b 1, Laptop 2, Mouse-2 3, Keyboard グループ A 1 2 4: D 3 INSERT グループ B (該当なし) id 4 は DELETE のみのため、 DELETE として後続に反映
  23. Leader Schedular ① IncrementalChangelogScan の 結果を WRITE タスクとして DB に登録

    ④ ③の完了後、DB からグループ情報を 抽出、空グループは除外・⼩さいグループ は束ねて READ タスクとして DB に登録 ⑥ ⑤の完了後、処理済 スナップショット id を記録 → 最初へ戻る ③ PK のハッシュでグループに ⑤ READ タスクを DB から取得し、 Node 1 グループサイズの変化に合わせてタスクを構築 ② WRITE タスクを DB から取得 分けてローカルディスクに Write 担当グループに対応するデータを Pull (COALESCE ) することで、 → メタデータを DB に登録 → carryover / UPDATE を認識 し、ファイルを Read グループ A データ流⼊量に応じて柔軟にタスク数を最適化できる グループ A 担当 DELETE file a 1 2 3 DELETE 1, Laptop 1: D+I 2: D+I 3: D+I 2, Mouse-1 グループ B id 1, 3 は全カラム同じ値なので相殺 3, Keyboard 4 DELETE id 2 は Mouse-1 がMouse-2 になっ 4, Monitor ているので UPDATE として後続反映 Node 2 INSERT file b 1, Laptop 2, Mouse-2 3, Keyboard Pull 型でデータを取得することで、リトライしやすい グループ B 担当 +ノード数 が変化しても、均等な分配を実現しやすい 4: D グループ A 1 2 3 INSERT グループ B (該当なし) id 4 は DELETE のみのため、 DELETE として後続に反映
  24. Spark Adaptive Query Execution の実装に着想を得た • Spark の shuffle はノード間でグルーピングしたデータを

    pull する仕組み • Spark AQE では、map タスクが shuffle 書き込みの完了後に、reduce partition のタスク統計情報を Driver が収集して、適切な reduce partition 数 に COALESCE して調整する • これによってユーザーが明⽰的に partition 数を指定しなくても、最適化される https://spark.apache.org/docs/latest/sql-performance-tuning.html
  25. パフォーマンステスト • ソース Iceberg テーブル︓NYC Yellow Taxi (4,100 万⾏、19 カラム、⽇別

    パーティション)→ シナリオごとに更新を実施 • Data Prepper ノードスペック︓ECS Fargate 2 vCPU / 16 GiB(ARM64) • ターゲット︓ OpenSearch: r7g.4xlarge × 6 • 測定範囲: スナップショット検出 → パイプライン Buffer 投⼊まで (polling 待ちと sink 側は除外) シナリオ Initial load (41M rows) CDC INSERT 100K CDC INSERT 1M CDC UPDATE 50K 8 ノード 2m11s 8s 10s 35s 16 ノード 1m26s 5s 6s 28s CDC DELETE 50K 37s 21s
  26. ロードマップ ≒ 今後やりたいこと • Merge on Read への対応 - Upstream

    の Iceberg へのコントリビュートとセットで要検討 - Row level Lineage の活⽤余地も考えたい • Coordination Store の拡張 #6740 - 現状 DynamoDB のみだが、PostgreSQL / MySQL への拡張を開発中 • 各種性能最適化 - 初期ロードの更なる⾼速化 #6725 - SHUFFLE_WRITE のグルーピングによる⾼速化 #6724 - SHUFFLE のオブジェクトストレージ等へのスピルの検討 - 複数スナップショットの読み⾶ばしによる⾼速化 #6667 36
  27. まとめ • Iceberg はオープンなテーブルフォーマットであり、テーブル仕様に準拠してい れば、オリジナルなツールからでも直接操作できる • Iceberg のリポジトリには、Iceberg を操作する上で便利なライブラリが公開さ れている

    • これらの仕組みを使って、 OpenSearch Data Prepper に「IcebergからOpenSearchへのCDC」を作った • IncrementalChangelogScan はスナップショット間の差分を把握するのに便利な ライブラリ • 設計を練る上で、Spark など先⼈の実装が⾮常に参考になった 37