Skip to content
Published on

ストリーミング vs バッチの再定義: Flink・RisingWave・Materialize、CDC、Streaming SQL、リアルタイムの実用主義 (2025)

シェア
Authors

Season 5 Ep 2 — Ep 1 が「データはどこに保存されるのか」だったとすれば、Ep 2 は「データはどれだけ速く流れるのか」。2020–2023 のリアルタイム狂騒から、2024–2025 の実用主義へ戻る。

Prologue — 「リアルタイムは基本ではなくオプションだ」

2019–2022年のデータカンファレンスのキーノートはどれも「すべてのデータはリアルタイムになるべきだ」だった。2025年の現実は違う:

  • リアルタイムパイプラインはバッチ比で 3–10倍高い
  • ほとんどの BI ダッシュボードは15分–1時間の遅延でも誰も文句を言わない
  • 本当にリアルタイムが必要なのは全データの 5–15%

2025年の正解:

「データごと・指標ごとの SLA を基準に鮮度階層を分け、各階層に合ったツールを使う」

この記事では、その階層とツール選定を具体化する。


第1章 · 鮮度階層(Freshness Tiers)

1.1 5階層のフレームワーク

階層レイテンシツール
Real-timems–秒取引モニタリング・異常検知・fraudFlink, Kafka Streams
Near-real-time1–5分運用ダッシュボード・alertFlink, RisingWave, Materialize
Fresh5–60分在庫・広告最適化Streaming append + rollup
Daily24時間BI, レポートSpark, dbt, SQL warehouse
Historical週/月分析・ML 学習Batch 年/月

1.2 各指標・テーブルを階層にマッピング

  • 在庫指標: Fresh
  • ログインの異常検知: Real-time
  • 月別売上: Daily
  • モデル学習用の特徴量: Fresh + Daily の混在

1.3 意思決定の原則

  • ビジネスの問い: 「指標が1時間遅れたら何が起きるのか?」
  • 答えが「何も起きない」→ Daily または Fresh
  • 答えが「10万ウォンの損失」→ Near-real-time
  • 答えが「人が怪我をする/法的問題になる」→ Real-time

第2章 · Lambda・Kappa アーキテクチャの2025年版

2.1 Lambda (2014)

  • Batch layer + Speed layer + Serving layer
  • 同じロジックを2回実装 → 維持の負担
  • 2015–2020 に支配的

2.2 Kappa (2014, Jay Kreps)

  • ストリーミングログ(Kafka)ひとつに統一
  • 再処理もストリーミングで
  • 理論的な単純さ、実務的な難しさ

2.3 2025年: Unified on Lakehouse

  • Iceberg/Delta テーブルの上にバッチとストリーミングが同時に書き込む
  • Flink + Iceberg, Spark Structured Streaming + Delta
  • Serving は OLAP エンジンまたは materialized view

2.4「Streaming + Materialized」パターン

  • ソース → ストリーミングパイプライン → Lakehouse
  • よく使うクエリは Materialized view でプリコンピュート
  • ad-hoc は OLAP エンジンが直接クエリ

第3章 · ストリーミングエンジン4大比較

  • 本物のイベント時間(Event time)ベース、watermark、exactly-once
  • 状態(Keyed state)が豊富、大規模プロダクションで検証済み
  • 学習曲線が急、運用が複雑
  • 2024–2025 で Flink CDC 2.x、Flink SQL が成熟

3.2 Spark Structured Streaming

  • Spark エコシステムの micro-batch ストリーミング
  • Databricks の標準
  • レイテンシ目標は 1–60秒以上、本物の ms 級ではない
  • SQL/Python が簡単

3.3 Kafka Streams / ksqlDB

  • Kafka 内蔵ライブラリ、JVM ベース
  • 軽量、Kafka のデータに集中するときに最適
  • 分散運用は自前で実装

3.4 RisingWave

  • 2022年オープンソース、Streaming SQL に特化
  • Postgres 互換 SQL + Materialized view
  • 状態は別ストレージ(S3)— 運用が簡単
  • Flink の代替として急浮上

3.5 Materialize

  • Streaming + Incremental view maintenance
  • 複雑なジョイン・集計を自動で最新化
  • エンタープライズ SaaS 中心

3.6 比較表

エンジンレイテンシ複雑度SQL韓国での利用特徴
Flinkms–秒高いO多い業界標準、状態管理が強い
Spark SS秒–分O非常に多いDatabricks と親和的
Kafka StreamsksqlDB普通Kafka 内蔵
RisingWave低いPostgres増加中運用が簡単、SaaS/OSS
MaterializeOIncremental view が強み

3.7 選定ガイド

  • イベント時間・複雑な状態: Flink
  • Databricks エコシステム: Spark SS
  • Kafka 中心・単純なロジック: Kafka Streams / ksqlDB
  • SQL だけでストリーミングアプリ: RisingWave または Materialize

第4章 · CDC (Change Data Capture)

4.1 なぜ CDC が核心なのか

  • 運用 DB(Postgres/MySQL)の変更を分析システムへリアルタイム複製
  • バッチダンプに比べて遅延・負荷が大幅に減少
  • Lakehouse の Silver レイヤーを自動的に維持

4.2 実装方式

  • Log-based: Postgres WAL, MySQL binlog — 最も効率的
  • Trigger-based: DB トリガーで変更イベントを生成 — DB 負荷が大きい
  • Query-based: updated_at カラムベースの polling — 簡単だが削除を捉えられない

4.3 ツール

  • Debezium (OSS): Postgres/MySQL/SQL Server/Oracle/Mongo → Kafka
  • Fivetran, Airbyte: マネージド ELT、数百のソースコネクタ
  • Striim, HVR: エンタープライズ CDC
  • Flink CDC: Debezium の上に Flink を統合

4.4 CDC → Iceberg パターン

PostgresDebeziumKafkaFlinkIceberg
  • Kafka でイベントを保管 → 再処理が可能
  • Flink が MERGE INTO で Iceberg に upsert
  • Iceberg v2/v3 の row-level delete を活用

4.5 実務の落とし穴

  • Schema change の処理(カラムの追加・削除・型変更)
  • 初期スナップショット + 増分(Incremental snapshot)の調和
  • Exactly-once の保証(重複・欠落の両方を防ぐ)
  • バックフィル・再処理の戦略

第5章 · Iceberg v3 とリアルタイム Upsert

5.1 Iceberg バージョンの歴史

  • v1: Append-only
  • v2: Row-level delete(position/equality)、MERGE INTO
  • v3 (2024–2025): Deletion vectors, row lineage, V3 partition transforms

5.2 Row-level delete の2方式

  • Position delete: 行のファイル・位置を指す削除ファイル
  • Equality delete: 条件ベースの削除(UPDATE の実装に有用)

5.3 リアルタイム Upsert のワークフロー

  1. Flink が CDC イベントを読む
  2. PK ベースで Equality delete + insert を生成
  3. Iceberg がスナップショットに反映
  4. 定期的な compaction で delete ファイルを整理

5.4 性能上の注意

  • Delete ファイルが溜まると読み取り性能が低下
  • Compaction の周期が重要(分–時間単位)
  • Copy-on-Write へ変換するバックグラウンド job の併用を推奨

第6章 · Streaming SQL の台頭

6.1 なぜ SQL なのか

  • Flink Java API は学ぶのが大変
  • データエンジニア + アナリストが一緒に扱える共通言語
  • dbt・Dagster のようなツールとの統合が容易
  • 2023–2024 に安定化
  • イベント時間・ウィンドウ・状態のすべてを SQL で
  • UDF・UDAF が可能

6.3 ksqlDB

  • Kafka の上の Streaming SQL
  • Table と Stream の概念が明確
  • 単純な ETL・集計に適する

6.4 RisingWave の Postgres 互換

  • CREATE MATERIALIZED VIEW → リアルタイムに維持
  • Postgres のツール(Grafana, Superset など)をそのまま利用
  • 運用の複雑さが低い → 2025年に急成長

6.5 Materialize

  • Incremental View Maintenance (IVM) の研究がベース
  • 複雑なジョインも自動で最新化
  • SaaS + セルフホストのオプション

第7章 · コストとレイテンシのトレードオフ

7.1 コストの構成

  • コンピュート(ストリーミングクラスタ 24/7)
  • 状態の保存(RocksDB/S3)
  • Kafka/MSK の運用
  • ネットワーク(リージョン間)

7.2 代表的なコスト比較(月額、中規模基準)

オプション月額コストレイテンシ
バッチ(Airflow + Spark、日次)低い($1–5k)24時間
Micro-batch(5分)中($3–10k)5分
Structured Streaming中–高($5–20k)秒–分
Flink クラスタ高($10–30k+)ms–秒
Managed(RisingWave/Confluent)中–高($7–25k)

7.3 削減テクニック

  • リアルタイム層を薄く(核心指標のみ)
  • Near-real-time 層で代替できないか検討
  • 閑散期のスケールダウン
  • 状態は S3(RisingWave 式)vs ローカル RocksDB(Flink)のトレードオフ

7.4 レイテンシ目標別の推奨アーキテクチャ

  • 100ms 未満: In-memory stream processor + Redis/RocksDB
  • 秒単位: Flink/RisingWave + Kafka
  • 分単位: Spark Structured Streaming + Delta/Iceberg
  • 時間単位: Micro-batch Airflow + dbt
  • 日単位: Batch Spark

第8章 · オブザーバビリティとデバッグ

8.1 主要指標

  • End-to-end latency(イベント発生 → 結果への反映)
  • Lag(Kafka consumer lag, Flink checkpoint lag)
  • Throughput(eps, rps)
  • Back-pressure(Flink task ごと)
  • 状態サイズ(RocksDB size, checkpoint size)

8.2 観測ツール

  • Flink UI + Prometheus + Grafana
  • Kafka Lag Exporter, Burrow, Conduktor
  • OpenLineage でストリーミングパイプラインの系譜
  • Datadog/NewRelic APM + ストリーミング extensions

8.3 デバッグ

  • チェックポイント・セーブポイントベースの再処理
  • CEP(Complex Event Processing)ルールのデバッグ
  • サンプリングログでイベントを追跡
  • テストは MiniCluster + embedded Kafka

8.4 アラート

  • Lag > しきい値 → Slack/PagerDuty
  • End-to-end の遅延 SLO 違反
  • チェックポイント失敗の連続発生

第9章 · 障害・復旧・SLA

9.1 SLA 設計

  • Availability: 99.9% = 月43分のダウンを許容
  • Freshness: イベント発生後 X 秒以内にクエリ可能
  • Correctness: 最終的な正確性(eventual)vs exactly-once

9.2 復旧戦略

  • Flink: チェックポイント(定期)、セーブポイント(手動)
  • Kafka: Replication factor 3, min.insync.replicas 2
  • Iceberg: スナップショット + ブランチでロールバック

9.3 再処理

  • Kafka retention を再処理期間の分だけ確保(7–30日が一般的)
  • または Iceberg の原本 → ストリーミング job を再実行
  • 再処理中の重複防止(冪等性の設計)

9.4 Multi-region

  • Kafka MirrorMaker 2 / Confluent Cluster Linking
  • Iceberg はストレージ複製(S3 Cross-region)
  • ストリーミングジョブは Active-Passive が一般的

第10章 · ストリーミング + Lakehouse の実践パターン

10.1 Medallion 上のストリーミング

  • Bronze: 原本イベントを append
  • Silver: 整形・重複排除・ジョインの upsert
  • Gold: 集計・指標

10.2 CDC → Silver

  • DB の変更イベント → Flink → Iceberg Silver に upsert
  • Silver は運用 DB の「コピー + 過去」の役割
  • Gold は dbt/Spark のバッチで生成

10.3 イベントソーシング

  • ドメインイベントを Kafka に永久保管
  • 状態はイベントの再生で再構成
  • Iceberg Bronze に長期保管

10.4 リアルタイム Feature Store

  • Feast/Tecton + Kafka + Iceberg
  • Online feature(Redis)+ Offline feature(Iceberg)
  • Online-offline skew のモニタリング

10.5 Streaming ETL パイプライン

  • ソース → Bronze(append)→ Silver(整形)→ Gold(集計)
  • 各段階を Flink/Spark SS のジョブで実装
  • 再処理はアップストリームの offset だけリセット

第11章 · 実践ケース3選

11.1 EC の注文パイプライン

  • 注文イベント Kafka → Flink → 在庫更新 + 異常検知 + 分析
  • 在庫: 秒単位のリアルタイム
  • 売上指標: 5分の near-real-time
  • 月次レポート: Daily batch

11.2 金融取引モニタリング

  • 取引イベント Kafka → Flink CEP → Fraud score
  • 遅延 100ms 未満を要求
  • 状態: ユーザーごとの取引履歴(Flink keyed state)
  • 結果: Redis + Iceberg に保存

11.3 ゲームテレメトリ

  • クライアントイベント Kinesis/Kafka → Flink/Spark SS
  • リアルタイム分析(同時接続者・DAU)は Near-real-time
  • 詳細ログは Iceberg Bronze に蓄積して後で分析
  • A/B テスト: Gold の集計

第12章 · 韓国企業のストリーミング

12.1 伝統的なパターン

  • 金融: Tibco EMS, IBM MQ, Kafka が混在
  • 通信: Charging/Billing にリアルタイムストリーミング
  • ゲーム: Kafka + Flink または Kafka Streams
  • コマース: Kafka + Spark SS + ELK

12.2 最新動向

  • Confluent Cloud / MSK の導入が増加
  • Flink の採用が拡大(特にトス・クーパン・ネイバー・カカオ)
  • RisingWave・Materialize は 2024–2025 が導入初期

12.3 規制上の考慮

  • 金融: 網分離環境で自前の Kafka クラスタ
  • 個人情報: CDC イベントの PII マスキングが必須
  • 監査ログ: 長期保管・不変性

12.4 難関

  • データエンジニアの人材不足 → マネージド志向が増加
  • レガシー DW との共存
  • 24/7 オンコール文化が定着途上

第13章 · アンチパターン10選

13.1「すべてをリアルタイムに」

必要のないテーブルまでストリーミング → コスト・複雑度が激増。

13.2 Exactly-once の盲信

ソース・シンクの両側で E2E exactly-once を保証するのは容易ではない。冪等設計が必須。

13.3 CDC の初期スナップショットを省略

欠落が発生し、正確性が低下。

13.4 スキーマ変更の自動伝播がない

ダウンストリームのパイプラインが壊れる。

13.5 Kafka retention が短すぎる

再処理が不可能になる。

13.6 チェックポイントの周期が長すぎる

障害時の復旧コスト・再処理量が激増。

13.7 状態を無制限に保管

Flink keyed state が無限に増加 → OOM。

13.8 Delete ファイルの compaction がない

Iceberg の読み取り性能が低下。

13.9 メモリ・CPU の過少割り当て

Back-pressure の連鎖。

13.10 オブザーバビリティ・アラートの不在

事故を顧客が先に発見する。


第14章 · チェックリスト — ストリーミングのローンチ前12項目

  • 鮮度階層のマッピング(テーブル・指標ごとの SLA)
  • エンジン選定の根拠(Flink/Spark SS/RisingWave/Materialize)
  • CDC ソースの選定 + 初期スナップショット戦略
  • Kafka retention + 再処理の計画
  • Iceberg/Delta テーブル設計 + compaction の自動化
  • チェックポイント・セーブポイントのポリシー
  • 冪等性・重複処理の設計
  • オブザーバビリティ(lag・latency・throughput・back-pressure)
  • SLA/SLO の定義 + アラート
  • コストダッシュボード(コンピュート・ストレージ・ネットワーク)
  • 災害復旧・Multi-region の計画
  • オンコール・障害対応のプレイブック

第15章 · 次回予告 — Season 5 Ep 3:「OLAP エンジン 2025 比較」

ストリーミングとバッチがストレージを共有するようになったなら、次の問いは「その上で誰が最も速くクエリするのか」だ。

  • DuckDB: シングルノード OLAP の革命
  • ClickHouse: リアルタイム OLAP の標準
  • Snowflake / BigQuery / Redshift: マネージドの巨人たち
  • Databricks SQL / StarRocks / Doris / Pinot / Druid
  • Trino / Presto の連合クエリ
  • 実際のベンチマークの罠
  • コスト vs レイテンシ vs 管理負担
  • MPP vs シングルノードの境界
  • 韓国企業の選定ガイド
  • 「適材適所」のエンジン配置パターン

1つのエンジンがすべてをこなすわけではない」という 2025年の現実を認めたあとが、本当に面白い。

次回の記事で会おう。


まとめ: 2025年のストリーミングは「すべてリアルタイム」から「SLA ベースの鮮度階層」へ再定義された。Real-time / Near-real-time / Fresh / Daily / Historical の5階層に合わせて Flink・Spark SS・RisingWave・Materialize・ksqlDB を配置し、CDC で運用 DB の変化を Lakehouse へ流し、Iceberg v3 の row-level delete でリアルタイム upsert を処理する。Lambda/Kappa ではなく「Unified on Lakehouse」が支配的なパターンであり、コスト・レイテンシ・複雑度のトレードオフを意識的に設計する。「リアルタイムは基本ではなくオプション」— これが 2025年の実用主義だ。