필사 모드: ストリーミング vs バッチの再定義: Flink・RisingWave・Materialize、CDC、Streaming SQL、リアルタイムの実用主義 (2025)
日本語Season 5 Ep 2 — Ep 1 が「データはどこに保存されるのか」だったとすれば、Ep 2 は「データはどれだけ速く流れるのか」。2020–2023 のリアルタイム狂騒から、2024–2025 の実用主義へ戻る。
- Prologue — 「リアルタイムは基本ではなくオプションだ」
- 第1章 · 鮮度階層(Freshness Tiers)
- 第2章 · Lambda・Kappa アーキテクチャの2025年版
- 第3章 · ストリーミングエンジン4大比較
- 第4章 · CDC (Change Data Capture)
- 第5章 · Iceberg v3 とリアルタイム Upsert
- 第6章 · Streaming SQL の台頭
- 第7章 · コストとレイテンシのトレードオフ
- 第8章 · オブザーバビリティとデバッグ
- 第9章 · 障害・復旧・SLA
- 第10章 · ストリーミング + Lakehouse の実践パターン
- 第11章 · 実践ケース3選
- 第12章 · 韓国企業のストリーミング
- 第13章 · アンチパターン10選
- 第14章 · チェックリスト — ストリーミングのローンチ前12項目
- 第15章 · 次回予告 — Season 5 Ep 3:「OLAP エンジン 2025 比較」
Prologue — 「リアルタイムは基本ではなくオプションだ」
2019–2022年のデータカンファレンスのキーノートはどれも「すべてのデータはリアルタイムになるべきだ」だった。2025年の現実は違う:
- リアルタイムパイプラインはバッチ比で 3–10倍高い
- ほとんどの BI ダッシュボードは15分–1時間の遅延でも誰も文句を言わない
- 本当にリアルタイムが必要なのは全データの 5–15%
2025年の正解:
「データごと・指標ごとの SLA を基準に鮮度階層を分け、各階層に合ったツールを使う」
この記事では、その階層とツール選定を具体化する。
第1章 · 鮮度階層(Freshness Tiers)
1.1 5階層のフレームワーク
| 階層 | レイテンシ | 例 | ツール |
|---|---|---|---|
| Real-time | ms–秒 | 取引モニタリング・異常検知・fraud | Flink, Kafka Streams |
| Near-real-time | 1–5分 | 運用ダッシュボード・alert | Flink, RisingWave, Materialize |
| Fresh | 5–60分 | 在庫・広告最適化 | Streaming append + rollup |
| Daily | 24時間 | 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大比較
3.1 Apache Flink
- 本物のイベント時間(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 | 韓国での利用 | 特徴 |
|---|---|---|---|---|---|
| Flink | ms–秒 | 高い | O | 多い | 業界標準、状態管理が強い |
| Spark SS | 秒–分 | 中 | O | 非常に多い | Databricks と親和的 |
| Kafka Streams | 秒 | 中 | ksqlDB | 普通 | Kafka 内蔵 |
| RisingWave | 秒 | 低い | Postgres | 増加中 | 運用が簡単、SaaS/OSS |
| Materialize | 秒 | 中 | O | 稀 | Incremental 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 パターン
Postgres → Debezium → Kafka → Flink → Iceberg
- 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 のワークフロー
- Flink が CDC イベントを読む
- PK ベースで Equality delete + insert を生成
- Iceberg がスナップショットに反映
- 定期的な compaction で delete ファイルを整理
5.4 性能上の注意
- Delete ファイルが溜まると読み取り性能が低下
- Compaction の周期が重要(分–時間単位)
- Copy-on-Write へ変換するバックグラウンド job の併用を推奨
第6章 · Streaming SQL の台頭
6.1 なぜ SQL なのか
- Flink Java API は学ぶのが大変
- データエンジニア + アナリストが一緒に扱える共通言語
- dbt・Dagster のようなツールとの統合が容易
6.2 Flink SQL
- 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年の実用主義だ。
현재 단락 (1/230)
2019–2022年のデータカンファレンスのキーノートはどれも「すべてのデータはリアルタイムになるべきだ」だった。2025年の現実は違う: