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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
3
总结2:KStream / KTable / GlobalKTable
| 抽象 | 本质 | 类比 | 使用场景 |
|---|---|---|---|
| KStream | 无界事件流 | 银行流水 | 用户行为、订单事件 |
| KTable | 变更日志 | 用户信息表 | 当前状态快照 |
| GlobalKTable | 全量副本表 | 字典/配置 | 维表 Join(数据量小) |
总结3:时间语义
Event Time: 业务时间(准确但需要处理乱序)
Processing Time: 处理时间(简单但不精确)
推荐默认使用 Event Time + Grace Period 处理迟到数据1
2
3
2
3
总结4:Kafka Connect
Source Connector: 外部系统 → Kafka(搬进来)
Sink Connector: Kafka → 外部系统(搬出去)
Debezium CDC: 最常用的 Source Connector,基于数据库 binlog1
2
3
2
3
总结5:技术选型
简单 SQL 流处理 → ksqlDB
复杂业务逻辑 → Kafka Streams (Java)
数据搬运 → Kafka Connect(零代码)1
2
3
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答案
答案:
| 维度 | KStream | KTable |
|---|---|---|
| 数据模型 | 插入流(每个事件独立) | 更新流(同 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:
- 配置 Debezium MySQL Connector,连接源库,指定要同步的表
- Debezium 实时读取 binlog,将 INSERT/UPDATE/DELETE 转换消息写入 Kafka Topic
- 配置 Elasticsearch Sink Connector,订阅对应的 Kafka Topic
- 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 的数据管道
学习状态:🟡 开始学习