Skip to content
Gains Summary
Main Navigation 首页 / Home
C++ 编程 / C++ Programming
系统与高性能 / Systems & Performance
Web 开发 / Web Development
人工智能 / Artificial Intelligence
工业软件 / Industrial Software
其他内容 / Other Topics
C++ 编程 / C++系统与性能 / SystemsWeb 开发 / Web人工智能 / AI工业软件 / Industrial

外观

Sidebar Navigation

← Web 开发 / Web Development

中间件 / Middleware

1. Web 中间件全景 / The Web Middleware Landscape

2. 中间件——流量洪峰下的系统保护 / Middleware for Protecting Systems Under Traffic Spikes

cache layer

1. 缓存架构全景 / Cache Architecture Overview

2. Redis 深入——为什么你的缓存总是出问题 / Redis Internals and Cache Failure Modes

3. Redis 数据结构场景应用——什么时候用什么 / Choosing Redis Data Structures for Real Applications

4. 缓存策略——库存变了,缓存怎么处理 / Cache Strategies and Inventory Consistency

5. Redis 缓存三剑客——穿透/击穿/雪崩 + 一致性策略 / Redis Cache Penetration, Breakdown, Avalanche, and Consistency

6. Redis 分布式锁——从 SETNX 到 Redisson / Redis Distributed Locks from SETNX to Redisson

7. Redis Cluster 与 Sentinel 高可用架构 / Redis Cluster and Sentinel High-Availability Architecture

8. Redis 高级特性——Stream / PubSub / Module / LLM 应用

9. Redis 架构深度分析——为什么 Redis 能这么快 / Redis Architecture and the Sources of Its Performance

10. Redis 全景——为什么你的系统需要一个缓存层 / The Redis Landscape and Why Systems Need a Cache Layer

message queue

1. 消息队列全景——为什么你的系统需要一个中间人 / Message Queue Overview and Why Your System Needs a Middleman

2. Kafka 核心——为什么你的消息总是"丢"了 / Kafka Fundamentals and Message Delivery Semantics

3. 消息队列高级——死信队列、延迟消息、消息积压 / Advanced Messaging with Dead Letters, Delays, and Backlogs

4. Kafka Streams 与 Connect —— 让数据自己流动起来 / Kafka Streams and Connect for Streaming Data Pipelines

5. RabbitMQ 深度解析 —— 灵活路由与消息可靠性 / RabbitMQ Deep Dive into Flexible Routing and Reliability

6. 消息队列对比——为什么最终选了 Kafka / Comparing Message Queues and Choosing Kafka

search engine

1. 搜索引擎知识体系 / Search Engine Knowledge System

2. Elasticsearch——为什么 Like 查询总是那么慢 / Elasticsearch for Full-Text Search at Scale

3. Elasticsearch 查询 DSL 深入——为什么你的搜索总是不准 / Elasticsearch Query DSL Deep Dive

4. Elasticsearch 集群规划与运维——为什么你的集群总是"黄" / Elasticsearch Cluster Planning and Operations

5. Meilisearch 与轻量搜索替代方案——当 ES 太重时 / Meilisearch and Lightweight Search Alternatives

infrastructure

1. 基础设施组件全景 / Infrastructure Components Landscape

2. Nginx 与反向代理 / Nginx and Reverse Proxy

3. 服务发现 / Service Discovery

4. 配置中心 / Configuration Center

本页目录

Kafka Streams 与 Connect —— 让数据自己流动起来 / Kafka Streams and Connect for Streaming Data Pipelines ​

📅 创建时间:2026-07-28 🏷️ 标签:#KafkaStreams #KafkaConnect #流处理 #ksqlDB #CDC #数据管道 📚 前置知识:[[/03-web/07-middleware/02-mq-kafka]](Kafka 核心机制) [[00-overview]](消息队列全景) 📚 相关知识:[[/03-web/07-middleware/07-mq-advanced]](MQ 高级特性) [[/03-web/06-databases-and-data-access/04-distributed-system]](分布式系统)


📋 本章目标 ​

  • 理解 Kafka Streams 的定位:流处理库,不是另一个集群
  • 掌握 KStream、KTable、GlobalKTable 的区别和使用场景
  • 理解流处理核心操作:map、filter、groupBy、join、windowing
  • 区分三种时间语义:Event Time、Processing Time、Ingestion Time
  • 理解 Kafka Connect 的架构:Source Connector -> Kafka -> Sink Connector
  • 了解常用 Connector:Debezium CDC、JDBC、S3、Elasticsearch
  • 初步认识 ksqlDB:用 SQL 写流处理

场景:为什么数据在 Kafka 里躺着,却没人用它 ​

┌─────────────────────────────────────────────────────────────┐
│                                                             │
│  你的 Kafka 集群里堆满了数据:                              │
│                                                             │
│  Topic "user-actions": 用户点击、浏览、下单                  │
│  Topic "orders": 订单数据                                   │
│  Topic "payments": 支付流水                                 │
│                                                             │
│  但问题是:                                                │
│  • 运营想看"每小时下单量排行" → 要写 Spark/Flink 任务     │
│  • 需要把 MySQL 增量数据同步到 Elasticsearch → 要写 Canal │
│  • 需要做实时风控(下单 3 次失败封号) → 要写流处理逻辑   │
│                                                             │
│  每个需求都要单独起一个项目,感觉 Kafka 只是一个"管子",  │
│  真正干活还得另找人。                                      │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17

Kafka 生态给出了答案:Kafka Streams 在 Kafka 上直接做流计算,Kafka Connect 在 Kafka 和外部系统间搬运数据。不需要 Spark 集群,不需要 Flink 集群,用 Kafka 就够了。


第1部分:Kafka Streams 是什么? ​

1.1 核心定位 ​

┌─────────────────────────────────────────────────────────────┐
│              Kafka Streams 的定位                              │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  Kafka Streams 是一个 Java 库,不是一个独立集群。            │
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  你的 Java 应用                                     │   │
│  │  ┌─────────────────────────────────────────────┐    │   │
│  │  │  Kafka Streams Library (JAR)                 │    │   │
│  │  │  • 从 Kafka 读数据                           │    │   │
│  │  │  • 在里面做计算(filter/map/join/aggregate) │    │   │
│  │  │  • 结果写回 Kafka                            │    │   │
│  │  └─────────────────────────────────────────────┘    │   │
│  │                                                     │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  对比:                                                     │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Spark/Flink:                                        │   │
│  │  ┌──────────┐    ┌──────────┐    ┌──────────┐      │   │
│  │  │  Kafka   │ →  │  Spark   │ →  │  Kafka   │      │   │
│  │  │ (Source) │    │ Cluster  │    │ (Sink)   │      │   │
│  │  └──────────┘    └──────────┘    └──────────┘      │   │
│  │       需要独立的计算集群!                             │   │
│  │                                                     │   │
│  │  Kafka Streams:                                      │   │
│  │  ┌──────────┐    ┌──────────────┐    ┌──────────┐  │   │
│  │  │  Kafka   │ →  │ 你的 Java 应用│ →  │  Kafka   │  │   │
│  │  │ (Source) │    │ (内含Streams) │    │ (Sink)   │  │   │
│  │  └──────────┘    └──────────────┘    └──────────┘  │   │
│  │       不需要独立集群!部署在你的应用里即可。           │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37

1.2 关键特性一览 ​

┌─────────────────────────────────────────────────────────────┐
│              Kafka Streams 的核心特性                          │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ✅ 精确一次语义(Exactly-Once Processing)               │
│     • 配合 Kafka 事务,处理结果不会重复                    │
│                                                             │
│  ✅ 有状态处理(Stateful Processing)                     │
│     • 本地 RocksDB 存储中间状态                            │
│     • 故障恢复时从 changelog Topic 重建状态                │
│                                                             │
│  ✅ 事件时间处理(Event-Time Processing)                 │
│     • 基于消息自带的时间戳,而非处理时间                    │
│     • 处理乱序数据和迟到数据                               │
│                                                             │
│  ✅ 弹性伸缩(Elastic Scaling)                           │
│     • 你的应用起几个实例,Streams 就自动分配 Partition    │
│     • 加实例 = 加算力                                     │
│                                                             │
│  ✅ 可组合(Composable)                                  │
│     • 多个处理步骤自然串联,代码清晰                       │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23

TypeScript/JavaScript 概念示例 —— 虽然不是 Java 但展示处理逻辑:

typescript
// Kafka Streams 的核心处理模式(概念代码,实际是 Java API)
// 源: Topic "user-actions" → 处理 → 汇: Topic "user-stats"

// 流处理拓扑 (Topology)
// Source → filter → map → groupBy → aggregate → to (Sink)

interface UserAction {
  userId: string;
  action: "view" | "click" | "purchase";
  productId: string;
  amount?: number;
  timestamp: number;  // Event Time
}

// 构建处理拓扑
function buildTopology(): void {
  const builder = new StreamsBuilder();

  // Step 1: 从源 Topic 创建 KStream
  const actions: KStream<string, UserAction> = builder.stream("user-actions");

  // Step 2: 过滤掉无效数据
  const validActions = actions.filter(
    (key, value) => value.userId != null && value.action != null
  );

  // Step 3: 只关心购买事件
  const purchases = validActions.filter(
    (key, value) => value.action === "purchase"
  );

  // Step 4: 按用户分组
  const byUser: KGroupedStream<string, number> = purchases
    .map((key, value) => KeyValue.pair(value.userId, value.amount ?? 0))
    .groupByKey();

  // Step 5: 聚合 → 每用户总消费金额(有状态操作)
  const userTotalSpent: KTable<string, number> = byUser.aggregate(
    () => 0,                                    // 初始值
    (userId, amount, total) => total + amount,  // 聚合逻辑
    Materialized.as("user-total-spent-store")   // 状态存储名
  );

  // Step 6: 结果写回 Kafka Topic
  userTotalSpent.toStream().to("user-total-spent-topic");

  // 启动
  const topology = builder.build();
  const streams = new KafkaStreams(topology, config);
  streams.start();
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51

第2部分:KStream vs KTable vs GlobalKTable ​

2.1 三种抽象的对比 ​

┌─────────────────────────────────────────────────────────────┐
│         KStream vs KTable vs GlobalKTable                     │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  KStream(事件流)—— 无界的插入流                           │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  每个事件都是独立的,流是无限的                      │   │
│  │                                                     │   │
│  │  user-actions stream:                               │   │
│  │  [login(u1)] → [view(u1)] → [click(u2)] → ...      │   │
│  │                                                     │   │
│  │  类比:银行流水(每笔交易独立记录,追加不停)        │   │
│  │  操作:filter、map、flatMap、foreach                │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  KTable(变更日志)—— 最新状态的快照                       │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  同一 Key 的新值覆盖旧值,始终反映最新状态            │   │
│  │                                                     │   │
│  │  user-profiles table:                               │   │
│  │  u1: {name: "张三"} → u1: {name: "张三丰"}          │   │
│  │                         (旧值被覆盖)                 │   │
│  │                                                     │   │
│  │  类比:用户信息表(只关心当前值,不关心历史)        │   │
│  │  操作:filter、mapValues、groupBy → aggregate      │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  GlobalKTable(全局表)—— 全量复制到每个实例               │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  数据量小,每个 Streams 实例都有完整副本              │   │
│  │                                                     │   │
│  │  country-codes table (200 条数据):                  │   │
│  │  CN → "中国", US → "美国", JP → "日本" ...         │   │
│  │                                                     │   │
│  │  类比:国家代码字典表(数据量小,全局共用)          │   │
│  │  用于和 KStream join,不需要重分区                  │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42

2.2 KStream 和 KTable 的 Join 语义 ​

┌─────────────────────────────────────────────────────────────┐
│          KStream-KTable Join 语义                             │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  KStream(订单流)Join KTable(用户信息表)                  │
│                                                             │
│  KStream: 订单事件(实时流入)                              │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  t1: {orderId: 1, userId: u1, amount: 100}          │   │
│  │  t2: {orderId: 2, userId: u2, amount: 200}          │   │
│  │  t3: {orderId: 3, userId: u1, amount: 150}          │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  KTable: 用户信息(当前最新状态)                            │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  u1: {name: "张三", vip: true}   ← t1 时是这样     │   │
│  │  u2: {name: "李四", vip: false}                     │   │
│  │  // t2.5: u1 变成 {name: "张三丰", vip: false}    │   │
│  │  u1: {name: "张三丰", vip: false}  ← t3 时是这样   │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  Join 结果(KStream-KTable Join):                         │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  t1: 订单1 + 张三 (VIP)     ← KTable 当时的快照     │   │
│  │  t2: 订单2 + 李四 (非VIP)                            │   │
│  │  t3: 订单3 + 张三丰 (非VIP) ← KTable 已更新         │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  关键:KStream 事件到达时才去查 KTable 的当前值             │
│  → 订单1 和 订单3 虽然是同一用户,但关联到的名字不同       │
│  → 这通常是正确的语义(订单关联下单时的用户信息)          │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
typescript
// 概念代码:KStream-Join-KTable
interface Order {
  orderId: string;
  userId: string;
  amount: number;
}

interface UserProfile {
  name: string;
  vipLevel: number;
}

interface EnrichedOrder {
  orderId: string;
  userId: string;
  userName: string;
  vipLevel: number;
  amount: number;
}

function enrichOrders(
  orders: KStream<string, Order>,
  userProfiles: KTable<string, UserProfile>
): KStream<string, EnrichedOrder> {
  return orders
    // 用 userId 做 Join Key
    .selectKey((key, order) => order.userId)
    .join(
      userProfiles,
      (order, profile) => ({
        orderId: order.orderId,
        userId: order.userId,
        userName: profile.name,
        vipLevel: profile.vipLevel,
        amount: order.amount,
      }),
      // Join 窗口:KTable 取值的时间范围
      Joined.with(Serdes.String(), orderSerde, profileSerde)
    );
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40

第3部分:流处理核心操作 ​

3.1 无状态操作 vs 有状态操作 ​

┌─────────────────────────────────────────────────────────────┐
│            无状态操作 vs 有状态操作                            │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  无状态操作(Stateless):                                  │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  每条消息独立处理,不依赖之前的消息                   │   │
│  │                                                     │   │
│  │  • filter:  msg → bool(通过/丢弃)                 │   │
│  │  • map:     msg → newMsg(转换)                    │   │
│  │  • flatMap: msg → [msg1, msg2, ...](一对多)       │   │
│  │  • branch:  msg → whichBranch(分流)               │   │
│  │  • peek:    msg → void(旁路观察,不改变流)        │   │
│  │                                                     │   │
│  │  特点:不需要本地存储,天然可并行                    │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  有状态操作(Stateful):                                   │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  需要"记住"之前见过的消息                             │   │
│  │                                                     │   │
│  │  • groupBy + count/aggregate/reduce: 分组聚合       │   │
│  │  • join: 关联两个流/表                              │   │
│  │  • windowing: 时间窗口聚合                          │   │
│  │                                                     │   │
│  │  特点:需要 RocksDB 本地存储状态                     │   │
│  │       故障恢复时从 changelog Topic 重建状态          │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30

3.2 Windowing 时间窗口 ​

┌─────────────────────────────────────────────────────────────┐
│              四种时间窗口类型                                  │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ① Tumbling Window(翻滚窗口)—— 固定大小,不重叠            │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  窗口大小: 1 分钟                                   │   │
│  │  [00:00-01:00]  [01:00-02:00]  [02:00-03:00]       │   │
│  │  每个窗口独立,无重叠                               │   │
│  │  适用:每分钟 PV 统计                               │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ② Hopping Window(滑动窗口)—— 固定大小,有重叠              │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  窗口大小: 5 分钟,滑动步长: 1 分钟                 │   │
│  │  [00:00-00:05]                                       │   │
│  │       [00:01-00:06]                                  │   │
│  │            [00:02-00:07]                             │   │
│  │  窗口之间有重叠                                     │   │
│  │  适用:近 5 分钟的热门商品(每分钟更新)             │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ③ Session Window(会话窗口)—— 基于活动间隔动态划分        │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  用户活动:  ██  ██        ██████    ██              │   │
│  │  会话:      [会话1]       [会话2]   [会话3]          │   │
│  │             (间隔<30min)                           │   │
│  │  适用:用户行为会话分析                             │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ④ Sliding Window(滑动窗口 Join 用)—— 基于时间差 Join     │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Join 窗口: 订单事件前后 1 小时内                   │   │
│  │  [订单t] → 找 [t-1h, t+1h] 内的支付事件            │   │
│  │  适用:订单-支付关联                               │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42

TypeScript 概念代码 —— 滚动窗口聚合:

typescript
// 每分钟统计各产品的购买量
function computePerMinutePurchaseStats(
  purchases: KStream<string, Purchase>
): void {
  purchases
    .groupBy(
      (key, purchase) => purchase.productId,
      Grouped.with(Serdes.String(), purchaseSerde)
    )
    // 1 分钟的翻滚窗口
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
    .count(Materialized.as("purchases-per-minute"))
    .toStream()
    // 将窗口时间转为可读格式
    .map((windowedKey, count) => {
      const productId = windowedKey.key();
      const windowStart = windowedKey.window().startTime();
      const windowEnd = windowedKey.window().endTime();
      return KeyValue.pair(
        productId,
        JSON.stringify({
          productId,
          window: `${windowStart} - ${windowEnd}`,
          count,
        })
      );
    })
    .to("product-stats-per-minute");
}
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29

第4部分:时间语义 ​

4.1 三种时间 ​

┌─────────────────────────────────────────────────────────────┐
│        Event Time vs Processing Time vs Ingestion Time        │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Event Time(事件时间)                              │   │
│  │  事件真正发生的时间(消息体里带的时间戳)            │   │
│  │  用户 12:00:05 点击 → 消息里 timestamp=12:00:05    │   │
│  │                                                     │   │
│  │  优点:反映业务真实时间                              │   │
│  │  挑战:网络延迟导致乱序                              │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Processing Time(处理时间)                         │   │
│  │  Streams 应用处理该消息时的系统时间                   │   │
│  │  应用在 12:03:20 处理 → processingTime=12:03:20     │   │
│  │                                                     │   │
│  │  优点:简单(不需要处理乱序)                        │   │
│  │  缺点:不反映业务真实情况                            │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Ingestion Time(摄入时间)                          │   │
│  │  消息到达 Kafka Broker 的时间                        │   │
│  │  Broker 在 12:01:30 收到 → ingestionTime=12:01:30  │   │
│  │                                                     │   │
│  │  介于 Event Time 和 Processing Time 之间            │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34

4.2 乱序和迟到数据处理 ​

┌─────────────────────────────────────────────────────────────┐
│              乱序数据处理策略                                  │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  问题:                                                     │
│  Event Time 窗口 [12:00, 12:05) 已经关闭了                  │
│  突然来了一条 Event Time = 12:02:30 的迟到消息              │
│  怎么处理?                                                │
│                                                             │
│  策略1:直接丢弃(Grace Period = 0)                        │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  窗口关闭后到达的消息直接丢弃                         │   │
│  │  适用:数据很规律,乱序容忍度低的场景                 │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  策略2:等待一段时间(Grace Period)                        │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  窗口关闭后,再等 N 分钟                             │   │
│  │  等待期内到达的迟到数据仍然更新窗口结果              │   │
│  │  等待期过后才真正关闭                               │   │
│  │                                                     │   │
│  │  窗口 [12:00, 12:05),grace = 1 分钟                │   │
│  │  12:05 窗口"关闭"但继续接受数据                     │   │
│  │  12:06 真正关闭并输出最终结果                        │   │
│  │                                                     │   │
│  │  代价:输出延迟增加了 grace period                   │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  策略3:迟到数据进旁路 Topic(Side Output)                  │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  Grace Period 过后仍到达的数据 → 写入专门的          │   │
│  │  "late-records" Topic,后续可做补算或人工处理        │   │
│  │                                                     │   │
│  │  适用:需要数据准确但对时效性要求不极端的场景        │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37

第5部分:Kafka Connect —— 数据搬运工 ​

5.1 Kafka Connect 是什么? ​

┌─────────────────────────────────────────────────────────────┐
│              Kafka Connect 架构                               │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  Kafka Connect 是一个框架,标准化地把数据搬进/搬出 Kafka。   │
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  ┌──────────┐    ┌──────────┐    ┌──────────┐      │   │
│  │  │  MySQL   │    │          │    │    S3    │      │   │
│  │  │ (Source) │──→│  Kafka   │──→│  (Sink)  │      │   │
│  │  └──────────┘    │          │    └──────────┘      │   │
│  │                  │          │                       │   │
│  │  ┌──────────┐    │          │    ┌──────────┐      │   │
│  │  │  PG/     │──→│          │──→│  Elastic- │      │   │
│  │  │  MongoDB │    │          │    │  search   │      │   │
│  │  └──────────┘    └──────────┘    └──────────┘      │   │
│  │                                                     │   │
│  │  Source Connector:  外部系统 → Kafka                │   │
│  │  Sink Connector:    Kafka → 外部系统                │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  核心价值:                                                 │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  ✅ 不需要写胶水代码                                 │   │
│  │  ✅ 统一管理(REST API 管理所有 Connector)          │   │
│  │  ✅ 容错(分布式部署,自动故障转移)                 │   │
│  │  ✅ 偏移管理(自动记录同步进度)                     │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31

5.2 Connect 配置示例 ​

typescript
// JDBC Source Connector —— 把 MySQL 数据同步到 Kafka
// 概念配置(实际是 JSON)
const jdbcSourceConfig = {
  "name": "mysql-orders-source",
  "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
  "tasks.max": "1",

  // 数据库连接
  "connection.url": "jdbc:mysql://localhost:3306/ecommerce",
  "connection.user": "connect_user",
  "connection.password": "${secret:mysql-password}",

  // 查询模式:增量轮询
  "mode": "timestamp",
  "timestamp.column.name": "updated_at",

  // 要同步的表 → 对应 Kafka Topic
  "table.whitelist": "orders",

  // Topic 命名规则
  "topic.prefix": "mysql-",

  // 轮询间隔
  "poll.interval.ms": "5000",
};

// Elasticsearch Sink Connector —— 把 Kafka 数据写入 ES
const esSinkConfig = {
  "name": "elasticsearch-orders-sink",
  "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
  "tasks.max": "2",

  // Kafka Topic → ES Index 映射
  "topics": "mysql-orders",
  "connection.url": "http://elasticsearch:9200",
  "type.name": "_doc",

  // 用消息 Key 做 ES 文档 ID(支持幂等写入)
  "key.ignore": "false",

  // 批量写入配置
  "batch.size": "1000",
  "linger.ms": "100",
};
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44

5.3 常用 Connector 一览 ​

┌─────────────────────────────────────────────────────────────┐
│              常用 Kafka Connector                              │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  Source Connectors(外部 → Kafka):                        │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  Debezium MySQL/PostgreSQL/MongoDB CDC              │   │
│  │  → 捕获数据库变更(INSERT/UPDATE/DELETE)           │   │
│  │  → 基于 binlog/WAL,实时且无侵入                    │   │
│  │  → 最常用的 CDC 方案,替代 Canal/Maxwell           │   │
│  │                                                     │   │
│  │  JDBC Source                                        │   │
│  │  → 轮询模式同步数据库表(基于 timestamp/自增ID)     │   │
│  │  → 比 CDC 简单,但有延迟                            │   │
│  │                                                     │   │
│  │  File/S3 Source                                     │   │
│  │  → 读取文件内容发送到 Kafka                         │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  Sink Connectors(Kafka → 外部):                           │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  JDBC Sink                                          │   │
│  │  → Kafka 消息写入关系型数据库                       │   │
│  │  → 支持 INSERT / UPSERT                             │   │
│  │                                                     │   │
│  │  Elasticsearch Sink                                 │   │
│  │  → Kafka 消息写入 ES(构建搜索索引)                │   │
│  │  → 支持批量写入、幂等(基于 Key)、动态 Index        │   │
│  │                                                     │   │
│  │  S3 Sink                                            │   │
│  │  → Kafka 消息写入 S3 对象存储                       │   │
│  │  → 按时间/大小分文件,适合数据湖场景                │   │
│  │  → 常用做长期存储和数据存档                         │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38

5.4 Debezium CDC 原理简述 ​

┌─────────────────────────────────────────────────────────────┐
│              Debezium CDC 工作流程                             │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  MySQL                                               │   │
│  │  ┌─────────────┐                                    │   │
│  │  │  Application │  INSERT/UPDATE/DELETE              │   │
│  │  └──────┬──────┘                                    │   │
│  │         ↓                                            │   │
│  │  ┌─────────────┐     ┌──────────────┐               │   │
│  │  │  binlog     │────→│ Debezium     │               │   │
│  │  │  (变更日志)  │     │ MySQL        │               │   │
│  │  └─────────────┘     │ Connector    │               │   │
│  │                      └──────┬───────┘               │   │
│  │                             ↓                       │   │
│  │                      ┌──────────────┐               │   │
│  │                      │    Kafka     │               │   │
│  │                      │  Topic:      │               │   │
│  │                      │  mysql.      │               │   │
│  │                      │  ecommerce.  │               │   │
│  │                      │  orders      │               │   │
│  │                      └──────┬───────┘               │   │
│  │                             ↓                       │   │
│  │               ┌────────────┴────────────┐           │   │
│  │               ↓            ↓            ↓           │   │
│  │          ES Sink     JDBC Sink    Streams 处理      │   │
│  │          (搜索)      (数仓)      (实时计算)        │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  每条变更消息包含:                                         │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  {                                                  │   │
│  │    "before": { "id": 1, "status": "pending" },     │   │
│  │    "after":  { "id": 1, "status": "paid" },        │   │
│  │    "source": { "table": "orders", "db": "ecom" },  │   │
│  │    "op": "u",        // c=create u=update d=delete │   │
│  │    "ts_ms": 1620000000000                           │   │
│  │  }                                                  │   │
│  │                                                     │   │
│  │  before + after → 可以做增量计算                    │   │
│  │  "查订单表的最新状态" 从查 MySQL 变成了消费 Kafka   │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46

第6部分:ksqlDB 简介 ​

6.1 什么是 ksqlDB? ​

┌─────────────────────────────────────────────────────────────┐
│              ksqlDB —— 用 SQL 写流处理                        │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ksqlDB 是建立在 Kafka Streams 之上的流式 SQL 引擎。        │
│  你写 SQL → ksqlDB 翻译成 Streams Topology → 在 Kafka 上跑  │
│                                                             │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  不用写 Java 代码,写 SQL 就行:                     │   │
│  │                                                     │   │
│  │  -- 创建流(从 Kafka Topic)                        │   │
│  │  CREATE STREAM orders (                             │   │
│  │    order_id VARCHAR KEY,                            │   │
│  │    user_id VARCHAR,                                 │   │
│  │    amount DOUBLE,                                   │   │
│  │    status VARCHAR                                   │   │
│  │  ) WITH (                                           │   │
│  │    KAFKA_TOPIC = 'orders-topic',                    │   │
│  │    VALUE_FORMAT = 'JSON'                            │   │
│  │  );                                                 │   │
│  │                                                     │   │
│  │  -- 过滤 + 持续查询                                 │   │
│  │  SELECT user_id, SUM(amount) AS total_spent         │   │
│  │  FROM orders                                        │   │
│  │  WHERE status = 'paid'                              │   │
│  │  GROUP BY user_id                                   │   │
│  │  EMIT CHANGES;                                      │   │
│  │                                                     │   │
│  │  -- 流-表 Join                                      │   │
│  │  CREATE STREAM enriched_orders AS                   │   │
│  │  SELECT o.order_id, o.amount, u.name, u.vip_level   │   │
│  │  FROM orders o                                      │   │
│  │  JOIN user_profiles u                               │   │
│  │  ON o.user_id = u.user_id                           │   │
│  │  EMIT CHANGES;                                      │   │
│  │                                                     │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  ksqlDB 适用场景:                                          │
│  ┌─────────────────────────────────────────────────────┐   │
│  │  ✅ 数据探索和原型开发                               │   │
│  │  ✅ 简单的 ETL pipeline                             │   │
│  │  ✅ 实时 Dashboard 数据源                            │   │
│  │  ✅ 数据质量监控(异常检测 SQL)                     │   │
│  │  ❌ 复杂业务逻辑(还是用 Kafka Streams Java)        │   │
│  │  ❌ 需要精确控制性能的场景                           │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50

6.2 ksqlDB 查询类型 ​

┌─────────────────────────────────────────────────────────────┐
│              ksqlDB 的两种查询模式                            │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  Pull Query(拉取查询)—— 查当前快照                       │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  SELECT * FROM user_total_spent                     │   │
│  │  WHERE user_id = 'u123';                           │   │
│  │                                                     │   │
│  │  → 返回当前这一刻 user_id=u123 的聚合结果           │   │
│  │  → 就像查数据库一样                                 │   │
│  │  → 底层从 RocksDB 状态存储中读取                    │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  Push Query(推送查询)—— 持续订阅变更                      │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  SELECT user_id, SUM(amount) AS total               │   │
│  │  FROM orders                                        │   │
│  │  GROUP BY user_id                                   │   │
│  │  EMIT CHANGES;                                      │   │
│  │                                                     │   │
│  │  → 结果持续更新,新数据来了就推送                    │   │
│  │  → 客户端收到增量变更,不是全量刷新                │   │
│  │  → 适合实时 Dashboard / 告警                        │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29

第7部分:实际架构模式 ​

7.1 常见数据管道组合 ​

┌─────────────────────────────────────────────────────────────┐
│           典型的 Kafka 数据管道组合                            │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  模式1:CDC → Kafka → 多目标同步(最常用)                   │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  MySQL ──→ Debezium ──→ Kafka ──→ ES Sink          │   │
│  │  (写入)   (binlog)      (中间)    (搜索索引)        │   │
│  │                           │                          │   │
│  │                           ├──→ S3 Sink(数据湖)     │   │
│  │                           ├──→ Redis Sink(缓存)    │   │
│  │                           └──→ Streams(实时聚合)   │   │
│  │                                                     │   │
│  │  解决:MySQL 一张表 → N 个异构系统数据同步问题      │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  模式2:实时 ETL(日志 → 清洗 → 聚合 → 存储)                │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  原始日志 → Kafka(raw) → Streams(清洗) → Kafka(clean)│   │
│  │                                     ↓                │   │
│  │                              Streams(聚合)           │   │
│  │                                     ↓                │   │
│  │                              Kafka(agg) → S3/ClickHouse│
│  │                                                     │   │
│  │  每 5 分钟聚合一次,存到 ClickHouse 做分析          │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
│  模式3:CQRS / Event Sourcing                               │
│  ┌─────────────────────────────────────────────────────┐   │
│  │                                                     │   │
│  │  写模型:Command → Kafka → 事件溯源                │   │
│  │  读模型:Kafka → Streams(聚合) → 物化视图          │   │
│  │                                                     │   │
│  │  "不要查数据库当前状态,消费 Kafka 构建自己的视图"  │   │
│  └─────────────────────────────────────────────────────┘   │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39

7.2 Streams 与 Connect 的配合 ​

┌─────────────────────────────────────────────────────────────┐
│         Streams + Connect 完整数据流                          │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  ┌─────────┐   ┌───────────┐   ┌──────────┐   ┌──────────┐ │
│  │  MySQL  │→  │ Debezium  │→  │          │   │          │ │
│  │ (订单表) │   │ Connector │   │          │   │          │ │
│  └─────────┘   └───────────┘   │          │   │          │ │
│                                │          │   │          │ │
│  ┌─────────┐   ┌───────────┐   │  Kafka   │   │  Kafka   │ │
│  │  App    │→  │ Producer  │→  │ (原始    │→  │ Streams  │ │
│  │ (日志)   │   │ SDK       │   │  数据)   │   │ (实时    │ │
│  └─────────┘   └───────────┘   │          │   │  计算)   │ │
│                                │          │   │          │ │
│                                └──────────┘   └────┬─────┘ │
│                                                    │       │
│                           ┌────────────────────────┘       │
│                           ↓                                │
│                    ┌──────────┐                            │
│                    │  Kafka   │                            │
│                    │ (结果    │                            │
│                    │  Topic)  │                            │
│                    └────┬─────┘                            │
│                         │                                  │
│              ┌──────────┼──────────┐                       │
│              ↓          ↓          ↓                       │
│        ┌─────────┐ ┌─────────┐ ┌─────────┐               │
│        │  ES     │ │  Redis  │ │  S3     │               │
│        │ Sink    │ │  Sink   │ │  Sink   │               │
│        └─────────┘ └─────────┘ └─────────┘               │
│                                                             │
│  Connect:负责数据搬运(外部 ↔ Kafka)                     │
│  Streams:负责数据计算(在 Kafka 内部做处理)               │
│  二者配合:Connect 搬运 → Kafka 存储 → Streams 计算        │
│                                                             │
└─────────────────────────────────────────────────────────────┘
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36

核心总结 ​

总结1:Kafka Streams 定位 ​

Kafka Streams = 嵌入你应用的流处理库(不是独立集群)
不需要 Spark / Flink 集群,一个 JAR 包搞定
通过增加应用实例数量实现弹性伸缩
1
2
3

总结2:KStream / KTable / GlobalKTable ​

抽象本质类比使用场景
KStream无界事件流银行流水用户行为、订单事件
KTable变更日志用户信息表当前状态快照
GlobalKTable全量副本表字典/配置维表 Join(数据量小)

总结3:时间语义 ​

Event Time: 业务时间(准确但需要处理乱序)
Processing Time: 处理时间(简单但不精确)
推荐默认使用 Event Time + Grace Period 处理迟到数据
1
2
3

总结4:Kafka Connect ​

Source Connector: 外部系统 → Kafka(搬进来)
Sink Connector:   Kafka → 外部系统(搬出去)
Debezium CDC: 最常用的 Source Connector,基于数据库 binlog
1
2
3

总结5:技术选型 ​

简单 SQL 流处理 → ksqlDB
复杂业务逻辑 → Kafka Streams (Java)
数据搬运 → Kafka Connect(零代码)
1
2
3

章节测试 ​

测试1:Kafka Streams 和新启动一个 Spark 集群的区别是什么? ​

测试2:KStream 和 KTable 的核心区别是什么? ​

测试3:Tumbling Window 和 Hopping Window 有什么区别? ​

测试4:Event Time 和 Processing Time 的区别是什么?Event Time 面临什么挑战? ​

测试5:Debezium 为什么比 JDBC Source Connector 更适合做 CDC? ​

测试6:以下场景分别适合用 Kafka Streams 还是 ksqlDB? ​

A. 实时订单数据清洗,需要复杂的 if-else 业务规则 B. 产品经理想自己写一个"每小时订单量统计" C. 需要 Exactly-Once 语义的支付流水聚合

测试7:如何构建一条从 MySQL 到 Elasticsearch 的实时数据同步管道? ​


参考答案 ​

测试1答案 ​

答案:

  • Kafka Streams 是一个 Java 库,嵌入你的应用进程运行,不需要独立部署集群
  • Spark 需要部署独立的计算集群(Master + Worker 节点)
  • Streams 实例数增加/减少时自动 Rebalance Partition,利用 Kafka 的 Consumer Group 机制
  • Streams 适合中等规模的流处理;Spark/Flink 适合需要独立计算资源的大规模场景

测试2答案 ​

答案:

维度KStreamKTable
数据模型插入流(每个事件独立)更新流(同 Key 新值覆盖旧值)
类比银行流水用户信息表
读取顺序消费所有事件读取当前最新快照
Join两个流的 Join 需要窗口Stream-Table Join 不需要窗口

测试3答案 ​

答案:

  • Tumbling Window:固定大小,窗口间不重叠(如每 1 分钟一个窗口)
  • Hopping Window:固定大小,但窗口间有重叠(如 5 分钟窗口每 1 分钟滑动一次)
  • Hopping 的每条消息可能属于多个窗口,计算量更大,但提供了更平滑的聚合结果

测试4答案 ​

答案:

  • Event Time:事件真实发生的时间(消息体中的时间戳),反映业务真实情况
  • Processing Time:Streams 应用处理该消息时的系统时间
  • Event Time 挑战:网络延迟、重试等导致消息乱序到达,需要通过 Grace Period 等待迟到数据,或使用 Side Output 处理超迟到数据
  • Event Time 是正确的水位线(Watermark)基础,是准确计算的必要条件

测试5答案 ​

答案:

  • Debezium 基于 MySQL binlog(WAL),实时捕获变更,延迟低(秒级),对源库无额外查询压力
  • JDBC Source 基于轮询(如按 timestamp 查询 WHERE updated_at > last_poll_time),有轮询间隔延迟,对源库有查询压力
  • Debezium 能捕获 DELETE 操作(JDBC 轮询无法感知删除)
  • Debezium 提供 before/after 完整变更记录,JDBC 只能看到最新状态

测试6答案 ​

答案:

  • A:Kafka Streams(复杂业务逻辑需要写代码)
  • B:ksqlDB(SQL 对非工程师友好,快速原型)
  • C:Kafka Streams(需要 Exactly-Once 语义的精确控制)

测试7答案 ​

答案: 使用 Debezium MySQL Source Connector(捕获 binlog)→ Kafka → Elasticsearch Sink Connector:

  1. 配置 Debezium MySQL Connector,连接源库,指定要同步的表
  2. Debezium 实时读取 binlog,将 INSERT/UPDATE/DELETE 转换消息写入 Kafka Topic
  3. 配置 Elasticsearch Sink Connector,订阅对应的 Kafka Topic
  4. Sink Connector 批量写入 ES,用消息 Key 做 ES 文档 ID 保证幂等

相关笔记 ​

  • [[/03-web/07-middleware/02-mq-kafka]] - Kafka 核心机制
  • [[/03-web/07-middleware/03-mq-others]] - 消息队列对比
  • [[00-overview]] - 消息队列全景
  • [[/03-web/06-databases-and-data-access/04-distributed-system]] - 分布式系统

下一步学习 ​

  • [ ] 阅读 [[04-rabbitmq-deep-dive]] - RabbitMQ 交换机、死信队列与高可用
  • [ ] 实践:用 Docker Compose 搭一套 Kafka + Debezium + ES 的数据管道

学习状态:🟡 开始学习

最后更新于:

Pager
上一篇3. 消息队列高级——死信队列、延迟消息、消息积压 / Advanced Messaging with Dead Letters, Delays, and Backlogs
下一篇5. RabbitMQ 深度解析 —— 灵活路由与消息可靠性 / RabbitMQ Deep Dive into Flexible Routing and Reliability

持续记录,持续成长

Copyright © Tidenflow