Chuyển đến nội dung chính

レッスン 15: ストリーム処理とリアルタイム データ パイプライン

バッチ処理とストリーム処理。 Apache Kafka ストリーム、基本的な Apache Flink。ウィンドウ戦略。変更データ キャプチャ (CDC)。リアルタイム分析パイプライン。ラムダ アーキテクチャとカッパ アーキテクチャ。

🏗️ アーキテクチャ — レッスン 15 レッスン 15: ストリーム処理とリアルタイム データ パイプライン

システムアーキテクチャ: ゼロからヒーローへ

パート 4: 非同期処理と通信

xdev.asia

はじめに

バッチ処理では、データをバッチ (時間ごと、日ごと) で処理します。しかし、リアルタイム時代では、ユーザーは結果をすぐに確認したいと考えています。 ストリーム処理 は、表示されたとおりにデータを処理します。


1. バッチ処理とストリーム処理

Batch Processing:
  ┌──────────┐    ┌──────────┐    ┌──────────┐
  │ Collect  │───►│ Process  │───►│ Output   │
  │ 24h data │    │ MapReduce│    │ Reports  │
  └──────────┘    └──────────┘    └──────────┘
  Latency: giờ → ngày
  Ví dụ: Daily sales report, ML training

Stream Processing:
  Event ──► Process ──► Output (liên tục)
  Event ──► Process ──► Output
  Event ──► Process ──► Output
  Latency: milliseconds → seconds
  Ví dụ: Fraud detection, live dashboard
基準バッチストリーム
レイテンシ分 → 時間さん → 秒
データ有界 (終点あり)無制限 (継続中)
処理データセット全体各イベント/マイクロバッチ
複雑さ下より高い
ツールスパーク、Hadoopカフカストリーム、フリンク

2. ストリーム処理の概念

2.1 ウィンドウ処理

Vấn đề: "Tính trung bình requests/giây trong 5 phút qua"
Cần nhóm events theo thời gian → Window

Tumbling Window (fixed, non-overlapping):
  |----5min----|----5min----|----5min----|
  | e1 e2 e3   | e4 e5 e6   | e7 e8      |
  | avg = 100  | avg = 150  | avg = 120  |

Sliding Window (overlapping):
  |----5min----|
       |----5min----|
            |----5min----|
  Events thuộc nhiều windows → smoother aggregation

Session Window (activity-based):
  |--user active--|  gap  |--user active--|
  | e1 e2 e3 e4   | 30m  | e5 e6          |
  | session 1     |      | session 2      |

2.2 イベント時間と処理時間

Event Time:     Khi event XẢY RA
Processing Time: Khi event ĐƯỢC XỬ LÝ

Vấn đề:
  Event tạo lúc 10:00:00 (event time)
  Network delay 5 giây
  Broker nhận lúc 10:00:05
  Consumer xử lý lúc 10:00:08 (processing time)

  Nếu dùng processing time → window 10:00-10:05 THIẾU event
  Nếu dùng event time → cần xử lý late events

Watermark:
  "Tôi tin rằng tất cả events trước thời điểm W đã đến"
  Events đến sau watermark → late events → xử lý riêng

3. Apache Kafka ストリーム

3.1 アーキテクチャ

  ┌──────────────┐      ┌──────────────┐
  │ Input Topic  │─────►│ Kafka Streams│─────► Output Topic
  │ "orders"     │      │ Application  │      "order-stats"
  └──────────────┘      │              │
                        │ Stateful     │
                        │ Processing   │
                        │ (RocksDB)    │
                        └──────────────┘

Đặc điểm:
  - Library (không phải cluster riêng)
  - Chạy trong Java/Kotlin app
  - State stores (RocksDB) cho aggregation
  - Exactly-once processing

3.2 例: リアルタイムの注文統計

StreamsBuilder builder = new StreamsBuilder();

KStream<String, Order> orders = builder.stream("orders");

// Đếm orders theo category trong 5 phút
KTable<Windowed<String>, Long> categoryCounts = orders
    .groupBy((key, order) -> order.getCategory())
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
    .count();

// Top products (số lượng bán)
KTable<String, Long> productCounts = orders
    .flatMapValues(order -> order.getItems())
    .groupBy((key, item) -> item.getProductId())
    .count();

categoryCounts.toStream().to("category-stats");
productCounts.toStream().to("product-rankings");

4. 変更データ キャプチャ (CDC)

4.1 CDC とは何ですか?

Vấn đề: Sync data giữa DB và Search/Cache/Analytics
  Option 1: Dual writes → Inconsistent (một cái fail)
  Option 2: Polling → Chậm, tốn resources
  Option 3: CDC → Capture DB changes → Stream

CDC captures database changes (INSERT, UPDATE, DELETE)
thành stream of events

  ┌──────────┐    ┌──────┐    ┌──────────────┐
  │PostgreSQL│───►│ CDC  │───►│ Kafka Topic  │
  │ WAL      │    │Debezium   │ "db.orders"  │
  └──────────┘    └──────┘    └──────┬───────┘
                                     │
                          ┌──────────┼──────────┐
                          ▼          ▼          ▼
                    Elasticsearch  Redis     Analytics
                    (Search index) (Cache)   (Data Lake)

4.2 Debezium CDC イベント

{
  "before": { "id": 1, "name": "Old Name", "status": "pending" },
  "after":  { "id": 1, "name": "New Name", "status": "shipped" },
  "source": {
    "connector": "postgresql",
    "db": "ecommerce",
    "table": "orders",
    "lsn": 12345678
  },
  "op": "u",
  "ts_ms": 1705312200000
}

4.3 使用例

1. Search Sync:
   DB → CDC → Kafka → Elasticsearch
   Mọi thay đổi DB tự động update search index

2. Cache Invalidation:
   DB → CDC → Kafka → Consumer → Redis.delete(key)
   Cache tự động invalidate khi data thay đổi

3. Cross-service Sync:
   Service A DB → CDC → Kafka → Service B
   Service B có read model của Service A data

4. Analytics:
   Production DB → CDC → Kafka → Data Lake
   Real-time data pipeline, không ảnh hưởng production

5. Lambda アーキテクチャと Kappa アーキテクチャ

5.1 ラムダアーキテクチャ

                    ┌─────────────────────────┐
                    │     Batch Layer          │
  Data ────────────►│  (Hadoop/Spark)          │
  Source    │       │  Complete, accurate      │──► Batch View
            │       │  High latency            │
            │       └─────────────────────────┘
            │
            │       ┌─────────────────────────┐
            └──────►│     Speed Layer          │
                    │  (Storm/Flink)           │──► Real-time View
                    │  Approximate, fast       │
                    │  Low latency             │
                    └─────────────────────────┘
                              │
                    ┌─────────▼───────────────┐
                    │    Serving Layer         │
                    │  Merge batch + realtime  │──► Query
                    └─────────────────────────┘

Nhược điểm: Maintain 2 pipelines (batch + speed)
            Duplicate logic

5.2 Kappa アーキテクチャ

  Data ──► Kafka (immutable log) ──► Stream Processing ──► View
  Source                              (Flink/Kafka Streams)

  Reprocessing:
    Replay Kafka log từ đầu
    → Rebuild views (thay vì batch job)

  Ưu điểm: 1 pipeline duy nhất
  Nhược điểm: Cần retained log đủ lâu
              Complex reprocessing

6. リアルタイム パイプラインの例

E-commerce Real-time Analytics:

  Web/App Events ──► Kafka "clickstream" ──┐
                                            │
  Order Events ──► Kafka "orders" ─────────┤
                                            │
  Payment Events ──► Kafka "payments" ─────┤
                                            ▼
                                    ┌──────────────┐
                                    │ Flink/Kafka  │
                                    │ Streams      │
                                    │              │
                                    │ • Sessionize │
                                    │ • Aggregate  │
                                    │ • Enrich     │
                                    │ • Detect     │
                                    │   anomalies  │
                                    └──────┬───────┘
                                           │
                              ┌────────────┼────────────┐
                              ▼            ▼            ▼
                        ┌──────────┐ ┌──────────┐ ┌──────────┐
                        │Dashboard │ │ Alert    │ │ Data     │
                        │(Grafana) │ │ System   │ │ Lake     │
                        │real-time │ │(PagerDuty│ │(long-term│
                        │metrics   │ │ Slack)   │ │analytics)│
                        └──────────┘ └──────────┘ └──────────┘

概要

コンセプト説明ツール
ストリーム処理リアルタイムイベント処理カフカストリーム、フリンク
ウィンドウ処理イベントを時間ごとにグループ化するタンブリング、スライディング、セッション
CDCDB の変更をキャプチャマクスウェル・デベジウム
ラムダバッチ + スピード レイヤーHadoop + ストーム
カッパストリーミングのみ、リプレイカフカ + フリンク

演習

  1. パイプライン設計: ソーシャル メディアにはリアルタイムのトレンド トピックが必要です。ストリーム処理パイプラインの設計: 入力 (投稿、いいね、共有) → 処理 → 出力 (毎分更新されるトレンド リスト)。

  2. CDC 実装: PostgreSQL (注文) と Elasticsearch (検索) があります。同期する CDC パイプラインを設計します。ケースの処理: Elasticsearch は 2 時間ダウンしましたが、その後回復しました。

  3. 窓口戦略: 不正検出: クレジット カードで 10 分間に 5 件を超える取引がある場合に警告します。どのような種類の窓を使用しますか?遅れたイベントにどう対処するか?