主题
第 17 章 Kafka Streams 与 ksqlDB:让 Kafka 自己处理流
目标读者:每次要做「实时统计 / Top-N / 实时大屏 / 实时风控」就拉一个 Flink 集群、写一堆 YAML,被 Flink 的部署和状态后端折腾得崩溃;或者干脆在 Consumer 里自己 hash table 聚合,结果重启后状态全丢的同学。
学完你会:能讲清「流-表二元性」、能在 5 分钟内写出 WordCount + 1 分钟滚动窗口 GMV 统计;能区分 KStream / KTable / GlobalKTable;能解释 Streams 怎么用 RocksDB + Changelog Topic 让状态既快又可恢复;能用 ksqlDB 写 SQL 跑流处理;能讲清楚 Streams 与 Flink 的本质区别。
0. 一个生活类比:电影 vs 截图
理解 Kafka Streams 之前,记住一个最重要的隐喻:
- 流(KStream) = 一部正在播放的电影(连续的事件序列:「每一帧画面」);
- 表(KTable) = 这部电影任意一刻按下暂停键的截图(按 key 聚合后的当前状态)。
电影完整记录了「每一帧」(不可改变的事实),截图就是某一帧的瞬时画面。两者可以无损转换:
- 电影 → 截图:从头到尾「重放」每一帧,最后一帧就是截图(聚合 / 物化);
- 截图 → 电影:把每张截图按时间排起来连播(Changelog)。
Kafka Streams 的核心创新就是:让流 (Topic) 和表 (State Store) 在你的代码里像同一种东西一样无缝转换。 这就是「流-表二元性」(Stream-Table Duality)。
1. Kafka Streams 是什么?为什么它不是另一个 Flink?
Kafka Streams 是一个 Java 库(不是集群、不是服务、不是框架),把它 import 进你自己的 Java 应用,应用启动 = 流处理任务启动;应用关停 = 任务停。无需 Flink JobManager / Spark Driver 这种「外部调度器」。
1.1 与其他流处理引擎对比
| 维度 | Kafka Streams | ksqlDB | Apache Flink | Spark Structured Streaming |
|---|---|---|---|---|
| 形态 | Java 库 | SQL 引擎 + REST | 集群 | 集群 |
| 部署 | 你自己 jar 部署 | docker / k8s | JobManager + TaskManager | Driver + Executors |
| 调度 | 无外部调度 | 自己集群 | 自带 | 自带 |
| 状态 | RocksDB + Changelog | 同 Streams | RocksDB + DFS Snapshot | RocksDB |
| 吞吐 | ★★★★ | ★★★ | ★★★★★ | ★★★★ |
| 延迟 | 毫秒 | 毫秒 | 毫秒 | 秒(micro-batch) |
| 上手 | Java DSL,门槛低 | SQL,最低 | 复杂 | 复杂 |
| 水平扩展 | 加实例(自动 Rebalance) | 加实例 | 加 TaskManager | 加 Executor |
| 与 Kafka 耦合 | 极强(仅 Kafka) | 极强(仅 Kafka) | 弱(多源) | 弱(多源) |
核心定位:Streams = 「最轻量级的流处理」,把流处理能力嵌入到普通微服务里;Flink 是「最强大的流处理集群」,独立运维。
一句话选型:
- 你只用 Kafka做事件源 + 想跑流处理 + 不想多维护一套集群 → Kafka Streams / ksqlDB;
- 你需要接 MySQL / Pulsar / 文件等多种 source、需要 CEP / 复杂窗口 / 跨源 join → Flink。
1.2 一段最简单的 Kafka Streams 代码
java
StreamsBuilder builder = new StreamsBuilder();
builder.<String, String>stream("learn.input")
.mapValues(v -> v.toUpperCase())
.to("learn.output");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();这就是一个完整的流处理应用:从 learn.input Topic 读 → 转大写 → 写到 learn.output。没有 JobManager、没有 Slot、没有 YAML。打成 jar 部署到 Kubernetes / Docker,加副本就是「水平扩展」。
2. 流-表二元性:Stream-Table Duality
这是整个 Streams 的灵魂。完整图示:
更精确的说法:
| 对象 | 数据语义 | 物理存储 | 等价对象 |
|---|---|---|---|
| KStream | 事件流:每条记录都是「发生了一件事」 | Kafka Topic | 数据库的 binlog |
| KTable | 当前状态:每个 key 的最新值 | 本地 RocksDB + Kafka Compacted Topic(Changelog) | 数据库的「当前表」 |
| GlobalKTable | 当前状态:每个 key 的最新值,全量复制到每个实例 | 同 KTable,但每个 instance 全量 | 维度表的「全表广播」 |
流和表的转换:
text
KStream.toTable() // 流 → 表(按 key 取最新值)
KTable.toStream() // 表 → 流(每次更新发一条变更)
KStream.groupByKey().count() // 流 → 聚合 → 表2.1 KStream / KTable / GlobalKTable 区别
| 维度 | KStream | KTable | GlobalKTable |
|---|---|---|---|
| 语义 | 事件序列 | 按 key 的当前值 | 按 key 的当前值,全量 |
| 重复 key | 各自独立事件 | 后值覆盖前值 | 同 KTable |
| 分区 | 按 key 分区 | 按 key 分区 | 每实例全量副本 |
| 容量 | 仅看 Topic 大小 | 单实例内存 + RocksDB | 必须全部塞进单实例 |
| 用途 | 业务事件、点击流 | 累加、维度的最新状态 | 小维度表(万级行) |
| Join 限制 | 与流 Join 必须共分区 | 与流 Join 必须共分区 | 不要求共分区(万能 Join 利器) |
记忆口诀:「事件用 KStream、状态用 KTable、广播维度表用 GlobalKTable」。
3. 拓扑(Topology):Streams 的「电路图」
Streams 程序在内部会被编译成一个 DAG(有向无环图),叫 Topology。Source 节点 → 算子节点 → Sink 节点。
可以用 streams.toString() 或调试模式打出来:
text
Topologies:
Sub-topology: 0
Source: source [topics: [learn.text]]
--> flatmap
Processor: flatmap (stores: [])
--> groupby
Processor: groupby (stores: [])
--> count
Processor: count (stores: [counts-store])
--> sink
Sink: sink (topic: learn.wc-output)3.1 Task 与 Partition
- 一个 Topology 可以拆成多个 Sub-Topology(用
repartition或groupBy切分); - 每个 Sub-Topology × 每个分区 = 一个 Stream Task;
- Task 是 Streams 的最小调度单元,一个 Task 由一个线程独占;
- 多个实例(Streams app 副本)通过 Consumer Group 协议 Rebalance 分配 Task。
text
Topic learn.text partitions=6 ⇒ 6 Stream Task
启 3 个实例,每个 num.stream.threads=2 ⇒ 6 个线程刚好一人一个 task
启 6 个实例,每个 num.stream.threads=1 ⇒ 一样 6 个 task,每实例 1 个水平扩展只需「加实例」,Rebalance 后任务自动迁移;状态 Store 也会通过 Changelog 在新实例上重建。
4. 无状态算子:流的基本工具箱
无状态 = 处理一条消息只看这条消息本身,不依赖任何累计。
| 算子 | 说明 | 示例 |
|---|---|---|
map | 同时改 key 和 value | (k,v) -> KeyValue.pair(k.toUpperCase(), v) |
mapValues | 只改 value(不破坏分区) | v -> v.toUpperCase() |
filter | 留下满足条件的 | (k,v) -> v.amount > 100 |
filterNot | 反向 filter | (k,v) -> v.status.equals("CANCELLED") |
flatMap | 一变多 | 拆 sentence → 词 |
flatMapValues | 一变多(仅 value) | s -> Arrays.asList(s.split(" ")) |
branch / split | 按条件分流到多个流 | 按金额分大单 / 小单 |
peek | 副作用(日志、metrics)不改流 | (k,v) -> log.info(...) |
foreach | 终止操作,仅副作用 | |
merge | 把两个 KStream 合成一个 | |
selectKey | 改 key(会触发 repartition) | (k,v) -> v.userId |
重要规则:改了 key 的算子(map / selectKey)会让流被打上 repartition 标记,下次有状态算子(groupBy / join)执行时,Streams 会偷偷在内部插一个 repartition Topic,把数据重新按新 key Hash 分区。这就是为什么有些应用启动后会发现 Kafka 多了几个 *-repartition Topic。
5. 有状态算子:流处理的「魂」
有状态 = 处理一条消息时要读 / 写一个累计的「状态」。背后必有 State Store + Changelog Topic 支撑。
5.1 groupByKey vs groupBy
java
// 已经按 key 分区,无需 repartition
KGroupedStream<String, Order> g1 = orders.groupByKey();
// 改了 key,会插 repartition
KGroupedStream<String, Order> g2 = orders.groupBy((k, v) -> v.userId);5.2 count / reduce / aggregate
最常见的三个聚合:
java
KGroupedStream<String, Order> grouped = orders.groupByKey();
// 1) count:每 key 计数 → KTable<String, Long>
KTable<String, Long> cnt = grouped.count();
// 2) reduce:把同 key 的 value 累计成一个相同类型的值
KTable<String, Order> latest = grouped.reduce((acc, cur) -> cur); // 取最新
// 3) aggregate:可以改输出类型;最强大、最常用
KTable<String, Stats> stats = grouped.aggregate(
() -> new Stats(0, 0.0), // initializer
(key, order, agg) -> agg.update(order.amount), // adder
Materialized.<String, Stats, KeyValueStore<Bytes,byte[]>>as("stats-store")
.withValueSerde(StatsSerde)
);5.3 join
| Join 类型 | 必须 window | 必须共分区 | 典型用例 |
|---|---|---|---|
KStream-KStream | ✅ 必须 | ✅ | 订单流 ⨝ 支付流(10 分钟内匹配) |
KStream-KTable | ❌ | ✅ | 订单流 ⨝ 用户表(流来一条查一次) |
KStream-GlobalKTable | ❌ | ❌ | 流 ⨝ 国家维度表(小表广播) |
KTable-KTable | ❌ | ✅ | 用户表 ⨝ 等级表(双表都按 key 维护) |
KStream-KStream join 示例:
java
KStream<String, Order> orders = ...;
KStream<String, Payment> payments = ...;
KStream<String, OrderWithPayment> joined = orders.join(
payments,
(o, p) -> new OrderWithPayment(o, p),
JoinWindows.of(Duration.ofMinutes(10))
);KStream-KTable join 示例:
java
KTable<String, User> users = builder.table("learn.users");
KStream<String, OrderEnriched> enriched = orders.join(
users,
(order, user) -> new OrderEnriched(order, user)
);5.4 窗口(Window):「用时间切流」
Streams 提供 4 种窗口:
| 窗口类型 | 时长 | 是否重叠 | 是否对齐 | 典型用例 |
|---|---|---|---|---|
| Tumbling(滚动) | 固定 | ❌ 不重叠 | ✅ 时间对齐 | 每分钟 GMV |
| Hopping(跳跃 / 滑动) | 固定 + step | ✅ 重叠 | ✅ 对齐 | 每 10s 输出最近 1 分钟 |
| Sliding(事件级滑动) | 固定 | 按事件 | ❌ | 反欺诈(最近 N 秒交易笔数) |
| Session(会话) | 不定 | ❌ | ❌ | 用户活跃 session(30min 无操作切断) |
生活类比:考勤表
- Tumbling = 「每天 9:00-18:00 的考勤」—— 严格不重叠;
- Hopping = 「每小时统计最近 4 小时考勤」—— 重叠;
- Session = 「员工进门到离门算一次会话」—— 自动伸缩;
- Sliding = 「每次有人刷卡就算最近 30 分钟的人数」—— 事件触发。
代码示例:
java
// 1 分钟 Tumbling
KTable<Windowed<String>, Long> perMinuteCnt = orders
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
.count();
// 1 分钟 Hopping,每 10 秒更新一次
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(5))
.advanceBy(Duration.ofSeconds(10)))
// 30 分钟 Session
.windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(30)))6. State Store + Changelog:状态可恢复的核心
每个有状态算子背后都有一个 State Store(默认 RocksDB),存的是「当前 key 的值」。
text
内存(RocksDB MemTable / BlockCache)
↕
本地磁盘(RocksDB SST 文件)
↕
Kafka Compacted Topic(Changelog,KAFKA-STREAMS-{app-id}-{store}-changelog)6.1 写路径
每次 State Store 更新(put / delete),Streams 会同步写一条 changelog Topic(compaction 模式),保证「磁盘 + Kafka」两份。
6.2 读路径
正常读:直接从 RocksDB 读,毫秒级。
6.3 故障恢复
实例宕机 / 迁移 → 新实例从 changelog Topic 重放,把 State Store 重建。这就是「状态可恢复」的本质:状态本身就是一个 compacted topic。
6.4 standby replicas
为加速故障恢复,可以配 num.standby.replicas=1:每个 Task 在另一个实例上保持一份热备的 State Store,正常时从 changelog 持续追,主实例挂掉后秒级接管。
properties
num.standby.replicas=1📌 与 Flink 对比:Flink 也用 RocksDB,但状态的 checkpoint 周期写到外部 DFS(HDFS / S3);Streams 没有 DFS 依赖,状态完全靠 Kafka Topic 持久化,运维更轻。
7. 时间语义:Event Time / Processing Time / Ingestion Time
| 时间 | 取自 | 适合 |
|---|---|---|
| Event Time | 消息 payload 里的业务时间戳(如 order.createdAt) | 业务真实时间,首选 |
| Processing Time | Streams 处理消息时的系统时间 | 简单的实时统计、不在意延迟 |
| Ingestion Time | Producer 把消息写进 Kafka 时由 broker 盖的时间戳 | 业务没有时间字段时的备选 |
7.1 自定义 TimestampExtractor
java
public class OrderEventTimeExtractor implements TimestampExtractor {
@Override
public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
Order o = (Order) record.value();
return o == null ? partitionTime : o.getCreatedAt();
}
}
// 注册
props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG,
OrderEventTimeExtractor.class);7.2 乱序与水位线
业务时间是不可避免会乱序的。Streams 用 stream time(截至此刻所有消息时间戳的最大值)当水位线,超过窗口 + grace period 的迟到消息会被丢弃。
java
TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofMinutes(5))Flink 用 Watermark,Streams 用 stream time + grace period,思路一致但 API 不同。
8. EOS v2:Streams 中的 Exactly Once
把第 13 章的事务 Producer + 幂等 Producer + read_committed 全部封装成一个开关:
java
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);启用后:
- 消费 source topic + 写 sink topic + 更新 State Store + 写 changelog 打包成一个事务;
- 提交一次性原子完成;
- 中途崩溃 → 整个事务 abort,下次重做;
- 下游消费者必须
isolation.level=read_committed才能看不到 abort 的中间状态。
EOS v2(Streams 2.6+)相比 v1,事务 Producer 数从「每 Task 1 个」减少到「每实例 1 个」,吞吐损耗从 ~30% 降到 ~10%。
代价:吞吐 / 延迟略降;commit interval(默认 100ms)越小延迟越低但开销越大。
9. 完整示例:WordCount
经典的 Hello World:
java
public class WordCount {
public static void main(String[] args) {
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
StreamsBuilder builder = new StreamsBuilder();
builder.<String, String>stream("learn.text")
.flatMapValues(v -> Arrays.asList(v.toLowerCase().split("\\W+")))
.filter((k, v) -> v != null && !v.isEmpty())
.groupBy((k, v) -> v)
.count(Materialized.as("counts-store"))
.toStream()
.to("learn.wc-output", Produced.with(Serdes.String(), Serdes.Long()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
streams.start();
}
}完整版(含 Maven 依赖、Serdes、运行说明)见 17_streams/code/wordcount_streams.java。
9.1 Python 等价方案:Faust
Kafka Streams 是 JVM 库,Python 没有官方等价。社区有 faust-streaming(前 Robinhood faust 的 fork),用 asyncio 风格提供类似 API:
python
import faust
app = faust.App('wordcount-py', broker='kafka://localhost:9092')
text_topic = app.topic('learn.text', value_type=str)
counts = app.Table('counts-store', default=int)
@app.agent(text_topic)
async def process(stream):
async for line in stream:
for word in line.lower().split():
counts[word] += 1
print(f'{word} -> {counts[word]}')Faust 不与 Streams 协议互通,状态后端是 RocksDB(基于
python-rocksdb),Changelog 也叫<app-id>-<table>-changelog但格式独立。生产环境 Faust 适合小规模 / Python 团队,不推荐替代 JVM Streams 跑高吞吐核心链路。
10. ksqlDB:用 SQL 写流处理
ksqlDB = Kafka Streams 的 SQL 外壳:你写 SQL,ksqlDB 把 SQL 编译成 Streams 拓扑,在 ksqlDB 集群里跑。
10.1 安装与启动(一行 Docker)
bash
docker run -d -p 8088:8088 \
-e KSQL_BOOTSTRAP_SERVERS=broker:9092 \
-e KSQL_LISTENERS=http://0.0.0.0:8088 \
-e KSQL_KSQL_SERVICE_ID=ksql-service \
confluentinc/ksqldb-server:0.29.010.2 创建 STREAM 和 TABLE
sql
-- 把 Kafka Topic 当 STREAM 暴露(KSTREAM)
CREATE STREAM clicks (
user_id VARCHAR KEY,
page VARCHAR,
ts BIGINT
) WITH (
KAFKA_TOPIC='learn.clicks',
VALUE_FORMAT='JSON',
TIMESTAMP='ts'
);
-- 把 Kafka Topic 当 TABLE(KTABLE,按 key 取最新值)
CREATE TABLE users (
user_id VARCHAR PRIMARY KEY,
name VARCHAR,
country VARCHAR
) WITH (
KAFKA_TOPIC='learn.users',
VALUE_FORMAT='JSON'
);10.3 流处理 SQL
sql
-- 1) 简单过滤 + 投影(持续查询)
CREATE STREAM tech_clicks AS
SELECT user_id, page
FROM clicks
WHERE page LIKE '/tech/%'
EMIT CHANGES;
-- 2) 1 分钟 Tumbling 窗口,按用户聚合
CREATE TABLE per_user_per_min AS
SELECT
user_id,
COUNT(*) AS click_cnt
FROM clicks
WINDOW TUMBLING (SIZE 1 MINUTE)
GROUP BY user_id
EMIT CHANGES;
-- 3) Stream-Table Join(流加维度表)
CREATE STREAM enriched_clicks AS
SELECT
c.user_id,
c.page,
u.country
FROM clicks c
LEFT JOIN users u
ON c.user_id = u.user_id
EMIT CHANGES;
-- 4) 30 分钟 Session 窗口(用户活跃 session)
SELECT
user_id,
COUNT(*) AS events,
WINDOWSTART AS session_start,
WINDOWEND AS session_end
FROM clicks
WINDOW SESSION (30 MINUTES)
GROUP BY user_id
EMIT CHANGES;EMIT CHANGES 表示「持续输出变更」,这是 ksqlDB 流处理的标志。EMIT FINAL 则只输出窗口关闭时的最终值。
10.4 ksqlDB 的优势
- 数据分析师 / SQL 用户上手 5 分钟;
- 拉起一个 docker 即可,无需写 Java;
- 默认 connector 集成(
CREATE SOURCE CONNECTOR ... WITH (...)); - pull query(
SELECT ... WHERE user_id='1')能像 KV 数据库一样查 KTable。
10.5 ksqlDB 的局限
- 不支持复杂的事件时间逻辑 / CEP;
- 没法写 UDAF(用户自定义聚合函数)以外的复杂状态机;
- 真正的「重业务流处理」还是要 Streams DSL 或 Flink。
11. 与 Flink 的最终对比
| 维度 | Kafka Streams | Apache Flink |
|---|---|---|
| 部署形态 | 普通 Java app | 独立集群 |
| 调度 | 无 | JobManager / TaskManager |
| Source | 仅 Kafka | 30+ Connector(Kafka / MySQL / 文件 / Pulsar / Kinesis…) |
| 状态后端 | RocksDB + Kafka Topic | RocksDB + DFS Snapshot |
| Exactly Once | EOS v2(仅 Kafka) | Two-Phase Commit(多 Sink) |
| CEP | ❌ | Flink CEP |
| 窗口 | Tumbling/Hopping/Sliding/Session | + Custom WindowAssigner |
| SQL | ksqlDB(独立) | Flink SQL(集成) |
| 学习曲线 | 平缓 | 陡 |
| 运维成本 | 低(无需独立集群) | 高 |
| 适合场景 | Kafka-only 微服务化流处理 | 多源、复杂、企业级 |
实战建议:
- 90% 的「实时统计 / 实时大屏 / 实时风控」用 Kafka Streams 或 ksqlDB 就够;
- 上 Flink 的理由通常是「需要接非 Kafka 数据源」、「跑 Flink CEP / SQL 已经成基础设施」、「跨多个团队共享同一个 Flink 集群」。
12. 生产踩坑
坑 1:忘了配 application.id 唯一
application.id 同时是 consumer group + 内部 changelog topic 前缀。两个不同的应用用同一个 application.id → state 互相覆盖,灾难现场。
坑 2:num.stream.threads 设过大
线程数 > 总 Task 数,多余线程空转浪费 RAM。配置时让 总实例数 × num.stream.threads ≈ Topic 总分区数 即可。
坑 3:State Store 撑爆磁盘
KTable / 窗口聚合的状态默认存 RocksDB(本地磁盘)。窗口太长 / key 基数太大 → 磁盘爆。 解决:
- 缩短 window;
- 调 RocksDB block cache(
rocksdb.config.setter); - 用
KStream+ 外部 KV(Redis)替代。
坑 4:repartition Topic 副本数 1
auto.create.topics.enable 为 true 时,Streams 内部创建的 repartition / changelog Topic 副本数 = 1,broker 重启可能丢数据。 解决:手动 streams.cleanUp() 后用 --replication-factor 3 创建,或配 replication.factor=3。
坑 5:EOS 开了但 commit.interval.ms 太大
EOS commit interval 决定可见性延迟(消息要等 commit 后下游才看得见)。默认 100ms,一些读者把它调成 30s 后抱怨「下游延迟太大」。
坑 6:Stream-Stream Join 没考虑 grace period
JoinWindows 默认 grace 24h(旧版)或 0(新版),新版默认 0 会丢迟到消息;旧版默认 24h 会让 State Store 累积大量历史。显式指定 grace 是好习惯。
坑 7:不同实例的状态规模不均
key 倾斜(如热门用户 ID)→ 一个 Task 状态特别大、其他空闲。需要在业务层散列 key(加随机后缀)或换聚合维度。
坑 8:直接重启实例改了 application.id
application.id 一变 = 全新应用 = 重新建 Changelog + 重读 source topic。线上千万别手抖改。
13. 小结
- Kafka Streams 是一个 Java 库(非集群),把流处理嵌进微服务,水平扩展靠加副本。
- 流-表二元性:KStream 是 binlog,KTable 是当前快照,二者可无损转换。
- 三种抽象:KStream(事件)/ KTable(状态)/ GlobalKTable(全量小维表)。
- 算子分两类:无状态(
map、filter、flatMap、branch、peek)和有状态(groupByKey、aggregate、reduce、count、join、window)。 - State Store = RocksDB + Changelog Topic(compacted),状态本质上是一个 Kafka Topic,因此「状态可恢复」是天生的。
- 三种 Join:KStream-KStream(必须 window)、KStream-KTable(流来一条查一次)、KTable-KTable(双表共维护)。
- 四种窗口:Tumbling / Hopping / Sliding / Session,配合时间语义(Event/Processing/Ingestion)+ grace period 处理乱序。
- EOS v2:一行
processing.guarantee=exactly_once_v2,把消费 + 写 + State 打包事务化。 - ksqlDB = Streams 的 SQL 外壳,用 SQL 描述流处理。
CREATE STREAM/TABLE+WINDOW TUMBLING即可。 - 与 Flink 选型:Kafka-only + 不想多维护集群 → Streams;多源 + 复杂 + 已有 Flink 基础设施 → Flink。
🎯 面试高频题
Q1:什么是流-表二元性?KStream / KTable / GlobalKTable 有什么区别?
考察点:核心概念、模型理解。
答案:
- 流-表二元性:流(KStream)是事件序列(每条都是「发生了一件事」),表(KTable)是按 key 的当前状态(最新值覆盖旧值)。流可以聚合成表(重放变更),表可以发出变更流(Changelog)。它们是同一信息的两种表达:流是电影,表是截图。
- KStream:每条消息独立,重复 key 不会覆盖;按 key 分区;适合事件、点击流、订单流。
- KTable:按 key 取最新值;按 key 分区;底层是 RocksDB + Changelog Compacted Topic;适合「当前状态」、累加结果、维度表。
- GlobalKTable:和 KTable 一样按 key 取最新值,但 每个实例都全量复制。适合小维度表(万级行),可以无视分区与流 join,是「广播 join」。
- 应用差别:
- 状态聚合 / 当前余额 → KTable;
- 事件 / 点击 → KStream;
- 国家代码、商品类目这种小维表 → GlobalKTable;
- 加分项:提到 KTable 与流 Join 必须共分区,GlobalKTable 不需要;提到 GlobalKTable 不参与 Stream Task 并行(每实例 1 份),所以不能太大。
Q2:State Store 是怎么做到「状态可恢复」的?
考察点:底层存储、容错。
答案:
- 本地存储:默认 RocksDB(key-value DB),存在每个 Streams 实例的本地磁盘;读写都是毫秒级。
- Changelog Topic:每次 State Store 更新(put / delete),Streams 同步写一条到
KAFKA-STREAMS-{app-id}-{store}-changelogTopic(compaction 模式),保证「磁盘 + Kafka」两份。 - 故障恢复:实例宕机 / 迁移 → 新实例从 Changelog 重放,把 RocksDB 重建。状态本质上就是一个 compacted topic。
- standby replicas:
num.standby.replicas=1让另一个实例热备一份 State Store,主挂秒级接管。 - Compacted Topic 的角色:保证「同 key 永远只保留最新值」,Changelog 体积可控;触发 Compaction 后旧值被回收。
- 对比 Flink:Flink 也用 RocksDB + checkpoint 到 DFS(HDFS / S3);Streams 没有 DFS 依赖,状态完全靠 Kafka,运维更轻。
- 加分项:提到 changelog 的 cleanup.policy=compact、副本数应 ≥ 2;提到 standby replicas 还能服务 Interactive Queries(IQ)。
Q3:四种窗口(Tumbling / Hopping / Sliding / Session)分别是什么?给典型场景。
考察点:窗口类型、应用场景。
答案:
- Tumbling(滚动):固定大小、不重叠、对齐时间。例:「每 1 分钟的 GMV」「每天的活跃用户数」。
- Hopping(跳跃 / 滑动):固定大小 + step;可重叠;对齐。例:「每 10 秒输出一次最近 1 分钟的实时 QPS 曲线」(窗口 1min,step 10s,相邻窗口重叠 50s)。
- Sliding(事件级滑动):每来一个事件就计算一次最近 N 时间内的状态;不对齐。例:「反欺诈 - 最近 30 秒同一用户交易笔数 > 5 报警」。
- Session:动态长度,由「无活动 gap」决定;相邻事件间隔超过 gap 就开新会话。例:「用户活跃 session(30 分钟无点击切断)」「设备上下线段统计」。
- 类比:考勤表里的不同打卡周期 —— Tumbling 像「每天 9:00-18:00」,Hopping 像「每小时统计最近 4 小时」,Session 像「员工进出门」,Sliding 像「每次打卡看最近 30 分钟」。
- 代码区别:java
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1))); .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)).advanceBy(Duration.ofSeconds(10))); .windowedBy(SlidingWindows.ofTimeDifferenceWithNoGrace(Duration.ofSeconds(30))); .windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(30))); - 加分项:提到 grace period 处理迟到消息;提到
EMIT CHANGES(每次更新)vsEMIT FINAL(窗口关闭再发)。
Q4:KStream-KStream Join、KStream-KTable Join、KTable-KTable Join 区别?
考察点:Join 模型、共分区。
答案:
- KStream-KStream:双流匹配,必须有 window(否则 State 无限增长)。例:订单流 ⨝ 支付流(10 分钟内匹配)。
- KStream-KTable:每来一条流消息查一次 KTable 当前值,不需要 window。例:订单流加用户档案。流来才触发 join,KTable 更新不触发。
- KTable-KTable:双表 join,不需要 window;任一边变化都触发 join。例:用户表 ⨝ 等级表。
- 共分区要求:以上三种都要求双方按相同 key 分区(不一致会自动 repartition,性能损耗)。唯一例外:与 GlobalKTable join 不要求共分区(小维表广播)。
- GlobalKTable 用途:当 join 的右侧是「小且全」(如国家代码、汇率表)时,用 GlobalKTable 避免 repartition + 全量在每实例本地查。
- 代码示例:java
// S-S window join ordersStream.join(paymentsStream, (o,p) -> ..., JoinWindows.of(Duration.ofMinutes(10))); // S-T join ordersStream.join(usersTable, (o,u) -> ...); // T-T join usersTable.join(levelsTable, (u,l) -> ...); - 加分项:提到 inner / left / outer join;提到 KStream-KTable 是 Streams 里最常用、性能最好的 join;提到 KTable-KTable foreign key join(2.4+)。
Q5:Kafka Streams 的 Exactly Once(EOS v2)是怎么实现的?跟 Flink 的 EOS 区别?
考察点:EOS 协议、流处理语义。
答案:
- 开关:
processing.guarantee=exactly_once_v2(Streams 2.6+)。 - 实现机制:
- 消费 source topic + 写 sink topic + 更新 State Store + 写 changelog 打包成一个事务 Producer 事务;
- 用
sendOffsetsToTransaction把 source consumer offset 也写到事务里; - commit 时一次性原子完成;
- 中途崩溃 → 整个事务 abort,下次从未提交的 offset 重做。
- 下游配合:消费者必须
isolation.level=read_committed,否则会读到中间未提交消息。 - v1 vs v2:v1 是「每个 Task 一个事务 Producer」,资源开销大、吞吐损耗 ~30%;v2 改成「每个实例一个事务 Producer」(Producer-per-instance),损耗降到 ~10%。
- 与 Flink EOS 区别:
- Streams:仅在 Kafka 内部生效(source、sink、state 都在 Kafka 上);
- Flink:基于 Two-Phase Commit Sink,可对接非 Kafka 系统(外部数据库);
- Streams 实现简单(依赖 Kafka 事务),Flink 实现通用(依赖 sink 的 2PC 协议)。
- 加分项:提到 commit.interval.ms 越小延迟越低但吞吐越低(默认 100ms);提到 EOS 不能解决「调用外部 API 副作用」的幂等(只在 Kafka topic + State 之间生效)。
Q6:Kafka Streams 是怎么水平扩展的?Task 和 Partition 是什么关系?
考察点:调度模型、Rebalance。
答案:
- 关系:
- 一个 Topology 被切成多个 Sub-Topology(按 repartition 划分);
- 每个 Sub-Topology × 每个分区 = 1 个 Stream Task;
- Task 是最小调度单元,一个 Task 由一个线程独占。
- 水平扩展:加 Streams 实例 → 通过 Consumer Group 协议 Rebalance → Task 自动迁移到新实例(含 State Store 重建 / 从 standby 接管)。
- 线程数:
num.stream.threads决定每实例多少线程;总线程 = 实例数 × 每实例线程数;线程数 ≤ Task 数才不浪费。 - 状态迁移:Task 迁到新实例时,新实例从 changelog 重放重建 State Store;如果有 standby replica,秒级接管。
- Rebalance 协议:Cooperative Sticky(默认 2.4+)—— 增量再平衡,避免 Stop-the-World。
- 类比 Flink:Flink 的 Slot ≈ Streams 的 Task;Flink 的 TaskManager ≈ Streams 实例;只是 Streams 没有 JobManager / 调度器,Rebalance 由 Kafka 协议自己搞定。
- 加分项:提到 Rebalance 听
KafkaStreams.setUncaughtExceptionHandler/StateRestoreListener监听状态恢复进度;提到分区数决定并行上限,Topic 分区数选少了 = 上限低。
本章配套:
17_streams/demo.html:流-表二元性 + 4 种窗口 + Stream-Table Join + WordCount Topology 可视化。17_streams/code/wordcount_streams.java:经典 WordCount Java 版。17_streams/code/wordcount_faust.py:Python Faust 等价实现。17_streams/code/realtime_topn.py:实时 Top-N 排行榜(Faust)。17_streams/code/windowed_aggregate.py:1 分钟滚动窗口聚合(Faust)。17_streams/code/ksql_examples.sql:ksqlDB 常见语句集合。
🔗 延伸阅读
- 第 11 章 消费者组与 Rebalance —— Streams 实例之间的 Cooperative Rebalance 是同一个协议。
- 第 13 章 幂等与事务 —— EOS v2 复用的事务 Producer + read_committed 协议。
- 第 14 章 Log Compaction —— Changelog Topic 用的就是 Compaction。
- 第 16 章 Kafka Connect —— Streams 处理后的结果可以用 Sink Connector 落地。
- 第 18 章 Schema Registry —— Streams 用 Avro Serdes 时依赖。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
sql
-- ====================================================================
-- 第 17 章 - Kafka Streams 与 ksqlDB
-- ksql_examples.sql - ksqlDB 常见语句集合
-- --------------------------------------------------------------------
-- 适用版本:ksqlDB 0.29.0+
-- 启动 ksqlDB CLI:
-- docker exec -it ksqldb-cli ksql http://ksqldb-server:8088
-- ====================================================================
-- ===== 0. 集群信息 =====
SHOW PROPERTIES;
SHOW STREAMS;
SHOW TABLES;
SHOW QUERIES;
SHOW TOPICS;
SHOW CONNECTORS;
-- ====================================================================
-- 1. 创建 STREAM(KSTREAM)= 把 Kafka Topic 当作事件流暴露
-- --------------------------------------------------------------------
-- 关键点:
-- * KEY 字段会作为消息 key
-- * TIMESTAMP 指定 event-time 来源(也可省略用 ROWTIME)
-- * VALUE_FORMAT:JSON / AVRO / PROTOBUF / DELIMITED
-- ====================================================================
CREATE STREAM clicks (
user_id VARCHAR KEY,
page VARCHAR,
referer VARCHAR,
ts BIGINT
) WITH (
KAFKA_TOPIC = 'learn.clicks',
VALUE_FORMAT = 'JSON',
TIMESTAMP = 'ts',
PARTITIONS = 6,
REPLICAS = 3
);
-- ====================================================================
-- 2. 创建 TABLE(KTABLE)= 按 key 取最新值
-- --------------------------------------------------------------------
-- 表的 source topic 必须是 compacted(或者你接受历史值被淘汰)
-- ====================================================================
CREATE TABLE users (
user_id VARCHAR PRIMARY KEY,
name VARCHAR,
country VARCHAR,
vip BOOLEAN
) WITH (
KAFKA_TOPIC = 'learn.users',
VALUE_FORMAT = 'JSON'
);
-- ====================================================================
-- 3. 持续查询(Push Query):流处理的灵魂
-- --------------------------------------------------------------------
-- EMIT CHANGES 表示「持续输出变更」,与传统 SQL 一次性返回不同
-- ====================================================================
-- 3.1 简单过滤 + 投影
CREATE STREAM tech_clicks AS
SELECT user_id, page
FROM clicks
WHERE page LIKE '/tech/%'
EMIT CHANGES;
-- 3.2 改 key(会触发内部 repartition)
CREATE STREAM clicks_by_page AS
SELECT page AS PAGE_KEY, user_id, ts
FROM clicks
PARTITION BY page
EMIT CHANGES;
-- ====================================================================
-- 4. 窗口聚合:4 种窗口
-- ====================================================================
-- 4.1 Tumbling 1 分钟:每分钟 PV
CREATE TABLE pv_per_minute AS
SELECT
page,
COUNT(*) AS pv,
COUNT_DISTINCT(user_id) AS uv
FROM clicks
WINDOW TUMBLING (SIZE 1 MINUTE)
GROUP BY page
EMIT CHANGES;
-- 4.2 Hopping:每 10s 输出一次最近 1min 的 QPS
CREATE TABLE pv_hopping AS
SELECT page, COUNT(*) AS pv_1m
FROM clicks
WINDOW HOPPING (SIZE 1 MINUTE, ADVANCE BY 10 SECONDS)
GROUP BY page
EMIT CHANGES;
-- 4.3 Session 窗口:用户活跃 session(30 分钟无操作切断)
CREATE TABLE user_sessions AS
SELECT
user_id,
COUNT(*) AS events,
TIMESTAMPTOSTRING(WINDOWSTART, 'HH:mm') AS session_start,
TIMESTAMPTOSTRING(WINDOWEND, 'HH:mm') AS session_end
FROM clicks
WINDOW SESSION (30 MINUTES)
GROUP BY user_id
EMIT CHANGES;
-- ====================================================================
-- 5. JOIN
-- ====================================================================
-- 5.1 Stream-Table Join:流 ⨝ 维度表(最常用、性能最好)
CREATE STREAM enriched_clicks AS
SELECT
c.user_id,
c.page,
u.name AS user_name,
u.country AS country,
u.vip AS is_vip
FROM clicks c
LEFT JOIN users u
ON c.user_id = u.user_id
EMIT CHANGES;
-- 5.2 Stream-Stream Join:双流匹配(必须 WITHIN)
-- 例:订单流 ⨝ 支付流(10 分钟内匹配)
CREATE STREAM orders_with_payment AS
SELECT
o.order_id,
o.amount,
p.tx_id AS payment_tx,
p.paid_at
FROM orders o
INNER JOIN payments p
WITHIN 10 MINUTES
ON o.order_id = p.order_id
EMIT CHANGES;
-- 5.3 Table-Table Join
CREATE TABLE user_with_level AS
SELECT
u.user_id,
u.name,
l.level
FROM users u
JOIN levels l
ON u.user_id = l.user_id
EMIT CHANGES;
-- ====================================================================
-- 6. Pull Query:像 KV 数据库一样查 KTable
-- --------------------------------------------------------------------
-- ⚠ 仅 TABLE 支持 pull query,且查询字段必须是 key
-- ====================================================================
SELECT * FROM users WHERE user_id = 'u-1';
-- 窗口聚合表:必须指定 WINDOW 范围
SELECT * FROM pv_per_minute
WHERE page = '/tech/article-1'
AND WINDOWSTART > 1700000000000;
-- ====================================================================
-- 7. CREATE STREAM AS SELECT 与持久查询管理
-- ====================================================================
-- 上面所有 CREATE STREAM/TABLE AS SELECT 都会生成一个「持久 Push Query」
SHOW QUERIES; -- 列出
EXPLAIN <query_id>; -- 看拓扑
TERMINATE <query_id>; -- 停止
DROP STREAM tech_clicks; -- 删除(必须先 TERMINATE 依赖)
-- ====================================================================
-- 8. UDF / UDAF(用户自定义函数)
-- --------------------------------------------------------------------
-- ksqlDB 内置丰富 UDF:UCASE / LCASE / SPLIT / REGEXP_EXTRACT / CASE WHEN
-- 时间:TIMESTAMPTOSTRING / STRINGTOTIMESTAMP / DATEADD / EXTRACT
-- JSON:EXTRACTJSONFIELD / JSON_OBJECT / TO_JSON_STRING
-- 集合:COLLECT_LIST / COLLECT_SET / EARLIEST_BY_OFFSET / LATEST_BY_OFFSET
-- --------------------------------------------------------------------
SELECT
user_id,
LATEST_BY_OFFSET(page) AS last_page,
COLLECT_LIST(page) AS visited_pages
FROM clicks
WINDOW TUMBLING (SIZE 5 MINUTES)
GROUP BY user_id
EMIT CHANGES;
-- ====================================================================
-- 9. Connector 集成(直接在 ksqlDB 里管理 Kafka Connect)
-- ====================================================================
CREATE SOURCE CONNECTOR debezium_orders WITH (
'connector.class' = 'io.debezium.connector.mysql.MySqlConnector',
'database.hostname' = 'mysql',
'database.port' = '3306',
'database.user' = 'debezium',
'database.password' = 'debezium-secret',
'database.server.id' = '1234',
'topic.prefix' = 'dbz.shop',
'table.include.list' = 'shop.orders'
);
SHOW CONNECTORS;
DESCRIBE CONNECTOR debezium_orders;
DROP CONNECTOR debezium_orders;
-- ====================================================================
-- 10. INSERT INTO(手动塞数据,调试方便)
-- ====================================================================
INSERT INTO clicks (user_id, page, ts)
VALUES ('u-1', '/tech/kafka', UNIX_TIMESTAMP());
-- ====================================================================
-- 11. SET 配置(仅本会话生效)
-- ====================================================================
-- 从最早开始消费
SET 'auto.offset.reset' = 'earliest';
-- 显式指定 EOS
SET 'processing.guarantee' = 'exactly_once_v2';
-- 增大并行度
SET 'num.stream.threads' = '4';
-- ====================================================================
-- 调试技巧:
-- * EXPLAIN <query_id>; 看 Streams Topology
-- * PRINT 'topic_name' FROM BEGINNING LIMIT 10; 打印 topic 内容
-- * DESCRIBE EXTENDED users; 看 schema、统计、所有依赖
-- ====================================================================python
"""
第 17 章 - Kafka Streams 与 ksqlDB
realtime_topn.py - 用 Faust 实现「实时 Top-N 商品销量排行榜」
业务场景:
电商大屏需要「最近 1 小时下单数 Top-10 商品」。每个订单是一个事件,
发到 `learn.orders` Topic;Streams 应用维护一个滑动窗口聚合 + Top-N
Heap,定时把 Top-10 推到 `learn.topn` Topic 给前端 BFF 拉取。
关键技巧:
1) 1 小时滚动窗口聚合每商品销量(KTable<product_id, count>)
2) 维护一个全局 Top-N Heap(Table,单 key='topn')
3) 每秒输出一次最新 Top-10 到下游 Topic / Web 接口
依赖:
pip install faust-streaming python-rocksdb
准备:
kafka-topics.sh ... --create --topic learn.orders --partitions 6 --replication-factor 1
kafka-topics.sh ... --create --topic learn.topn --partitions 1 --replication-factor 1
运行:
faust -A realtime_topn worker -l info --web-port 6068
curl http://localhost:6068/topn/now
"""
from __future__ import annotations
import faust
import heapq
from datetime import timedelta
from typing import List, Tuple
WINDOW_SIZE = timedelta(hours=1)
EMIT_EVERY = timedelta(seconds=5)
TOP_N = 10
class Order(faust.Record, serializer="json"):
order_id: str
product_id: str
user_id: str
amount: float
ts: float # event time, seconds
class TopNEntry(faust.Record, serializer="json"):
product_id: str
count: int
app = faust.App(
"realtime-topn",
broker="kafka://localhost:9092",
store="rocksdb://",
table_standby_replicas=1,
)
orders_topic = app.topic("learn.orders", value_type=Order)
topn_topic = app.topic("learn.topn", value_type=List[TopNEntry])
# 1 小时滚动窗口聚合:每个 product_id 的下单笔数
hourly_count = (
app.Table("hourly-count", default=int, partitions=6)
.tumbling(WINDOW_SIZE, expires=timedelta(hours=2))
.relative_to_field(Order.ts)
)
# 用一个 single-key Table 存当前 Top-N 快照(方便 IQ 查询)
topn_snapshot = app.Table("topn-snapshot", default=list, partitions=1)
# -------------------------------------------------------------------
# Agent 1:消费订单流,更新窗口聚合
# -------------------------------------------------------------------
@app.agent(orders_topic)
async def aggregate(stream):
"""每来一个订单:对应 product_id 的窗口计数 +1。"""
async for order in stream:
# Faust 的 windowed Table:直接用 [] 操作即可
hourly_count[order.product_id] += 1
# -------------------------------------------------------------------
# Timer:定期扫描全表,算出 Top-N,写到 snapshot + Kafka
# -------------------------------------------------------------------
@app.timer(interval=EMIT_EVERY.total_seconds())
async def emit_topn():
"""周期性把当前窗口的 Top-N 推送到下游 + 内存快照。"""
# 当前窗口下的所有 (product, count)
# 注意:Faust 的 windowed Table 遍历需要 .items() + .current()
pairs: List[Tuple[str, int]] = []
for k, w in hourly_count.items():
try:
cnt = w.current()
if cnt > 0:
pairs.append((k, cnt))
except Exception:
continue
# heapq.nlargest = Top-N
top = heapq.nlargest(TOP_N, pairs, key=lambda x: x[1])
payload = [TopNEntry(product_id=p, count=c) for p, c in top]
# 写到 snapshot Table(IQ 查询)
topn_snapshot["current"] = [(e.product_id, e.count) for e in payload]
# 写到下游 Topic
await topn_topic.send(value=payload)
if payload:
print(f"[Top-{TOP_N}] " +
" | ".join(f"{e.product_id}={e.count}" for e in payload))
# -------------------------------------------------------------------
# Web 接口:实时返回 Top-N(Interactive Query)
# -------------------------------------------------------------------
@app.page("/topn/now")
async def http_topn(web, request):
items = topn_snapshot.get("current", [])
return web.json([{"product_id": p, "count": c} for p, c in items])
if __name__ == "__main__":
app.main()
# ============================================================================
# 性能与扩展
# ----------------------------------------------------------------------------
# - 商品总数 1M+ 时,timer 里全表扫描会变慢,建议改成「维护 Top-N 增量更新」
# (比如只在 cnt 变化时检查是否能替换 heap 末尾)。
# - 窗口越长,State Store 越大;可以加 `compacting` 后端 + 缩短 expires。
# - 真正高吞吐(百万 QPS)建议用 Java Streams 或 Flink;Faust 里 Python
# 单进程吞吐有 GIL 上限。
# ============================================================================python
"""
第 17 章 - Kafka Streams 与 ksqlDB
windowed_aggregate.py - 用 Faust 演示「1 分钟滚动窗口聚合 GMV」
业务场景:
实时大屏需要「按分钟聚合的订单总额(GMV)」,按支付状态分组,
输出到下游 `learn.gmv-per-minute` Topic 给前端拉取。
核心知识点:
1) Tumbling Window(1 分钟,不重叠)
2) Aggregate(不仅 count,还累加 amount)
3) Event Time(用消息里的 ts 而不是处理时间)
4) Grace Period 处理乱序消息
依赖:
pip install faust-streaming python-rocksdb
准备:
kafka-topics.sh ... --create --topic learn.orders --partitions 3
kafka-topics.sh ... --create --topic learn.gmv-per-minute --partitions 1
运行:
faust -A windowed_aggregate worker -l info --web-port 6069
"""
from __future__ import annotations
import faust
from datetime import timedelta
from typing import NamedTuple
WINDOW_SIZE = timedelta(minutes=1) # 1 分钟滚动
GRACE_PERIOD = timedelta(seconds=30) # 允许 30s 迟到
class Order(faust.Record, serializer="json"):
order_id: str
user_id: str
amount: float
status: str # PENDING / PAID / CANCELLED / REFUNDED
ts: float # event time (epoch seconds)
class MinuteStats(faust.Record, serializer="json"):
bucket_minute: int # epoch minute
status: str
order_cnt: int
gmv: float
avg_amount: float
app = faust.App(
"windowed-gmv",
broker="kafka://localhost:9092",
store="rocksdb://",
)
orders_topic = app.topic("learn.orders", value_type=Order)
gmv_topic = app.topic("learn.gmv-per-minute", value_type=MinuteStats)
# -------------------------------------------------------------------
# 一个 windowed Table,key = status,value = (count, sum)
# -------------------------------------------------------------------
class Acc(NamedTuple):
cnt: int
sum: float
minute_acc = (
app.Table("minute-acc",
default=lambda: Acc(0, 0.0),
partitions=3)
.tumbling(WINDOW_SIZE, expires=timedelta(hours=1))
.relative_to_field(Order.ts) # 按 event time 做窗口
)
# -------------------------------------------------------------------
# Agent:每来一条订单,按 status 累加 cnt + sum
# -------------------------------------------------------------------
@app.agent(orders_topic)
async def consume_orders(stream):
async for order in stream:
cur = minute_acc[order.status].value()
minute_acc[order.status] = Acc(cur.cnt + 1, cur.sum + order.amount)
# -------------------------------------------------------------------
# 周期性扫描当前窗口,输出聚合
# -------------------------------------------------------------------
@app.timer(interval=10.0)
async def flush_to_topic():
"""每 10 秒把当前窗口的快照推到下游 Topic。"""
for status, w in minute_acc.items():
try:
acc: Acc = w.current()
except Exception:
continue
if acc.cnt == 0:
continue
# current() 返回当前 wall-clock 落入的窗口聚合
# 拿当前的 minute bucket
import time
bucket = int(time.time() // 60)
stats = MinuteStats(
bucket_minute=bucket,
status=status,
order_cnt=acc.cnt,
gmv=round(acc.sum, 2),
avg_amount=round(acc.sum / acc.cnt, 2),
)
await gmv_topic.send(key=f"{bucket}:{status}", value=stats)
print(f"[GMV] minute={bucket} status={status} "
f"cnt={stats.order_cnt} gmv={stats.gmv} avg={stats.avg_amount}")
# -------------------------------------------------------------------
# 实时查询接口:当前窗口的实时数字
# -------------------------------------------------------------------
@app.page("/gmv/current")
async def http_current(web, request):
out = []
for status, w in minute_acc.items():
try:
acc = w.current()
if acc.cnt > 0:
out.append({"status": status, "cnt": acc.cnt,
"gmv": round(acc.sum, 2)})
except Exception:
pass
return web.json(out)
if __name__ == "__main__":
app.main()
# ============================================================================
# 等价的 Java Streams DSL 写法(仅参考,不在本文件中运行)
# ----------------------------------------------------------------------------
# StreamsBuilder builder = new StreamsBuilder();
# KStream<String, Order> orders = builder.stream("learn.orders",
# Consumed.with(Serdes.String(), orderSerde)
# .withTimestampExtractor(new OrderTsExtractor()));
#
# KTable<Windowed<String>, Acc> agg = orders
# .groupBy((k, v) -> v.status)
# .windowedBy(TimeWindows.ofSizeAndGrace(
# Duration.ofMinutes(1),
# Duration.ofSeconds(30)))
# .aggregate(
# () -> new Acc(0, 0.0),
# (k, order, acc) -> new Acc(acc.cnt + 1, acc.sum + order.amount),
# Materialized.<String, Acc, WindowStore<Bytes,byte[]>>as("minute-agg")
# .withValueSerde(accSerde));
#
# agg.toStream()
# .map((wk, acc) -> KeyValue.pair(
# wk.window().start() / 60_000 + ":" + wk.key(),
# new MinuteStats(...)))
# .to("learn.gmv-per-minute", Produced.with(Serdes.String(), statsSerde));
# ============================================================================python
"""
第 17 章 - Kafka Streams 与 ksqlDB
wordcount_faust.py - 用 Python `faust-streaming` 实现的 WordCount
⚠ 注意:
Kafka Streams 是 JVM 库,Python 没有官方等价物。`faust-streaming` 是社区在
Robinhood `faust` 基础上 fork 出来的活跃版本,提供 Streams 风格的 API(agent +
Table),底层用 RocksDB + Changelog Topic。
Faust 的 Changelog 命名是 `<app-id>-<table>-changelog`,但格式与 Java Streams
并不互通;不能在生产中混用。
Faust 适合:Python 团队、轻量场景、原型快速验证;不推荐替代 JVM Streams
跑高吞吐核心链路。
依赖:
pip install faust-streaming python-rocksdb
准备 Kafka Topic:
kafka-topics.sh --bootstrap-server localhost:9092 \\
--create --topic learn.text --partitions 3 --replication-factor 1
运行:
# Worker 进程(可启多个,自动按分区 Rebalance)
faust -A wordcount_faust worker -l info --web-port 6066
# 喂数据
kafka-console-producer.sh --bootstrap-server localhost:9092 --topic learn.text
> kafka streams kafka topic kafka
# 查看 Table(实时排行)
curl http://localhost:6066/wordcount/top
"""
from __future__ import annotations
import faust
import re
from collections import OrderedDict
# ---------- App / Topic / Table ----------
app = faust.App(
"wordcount-py", # 等价 application.id
broker="kafka://localhost:9092",
value_serializer="raw", # 不要尝试 json 解码原始字符串
store="rocksdb://", # State Store 后端
table_standby_replicas=1, # 等价 num.standby.replicas
topic_partitions=3,
consumer_auto_offset_reset="earliest",
)
text_topic = app.topic("learn.text", value_type=bytes)
# Table 等价 KTable,落地到 RocksDB + Changelog Topic
counts = app.Table("counts-store", default=int, partitions=3)
WORD_RE = re.compile(r"\W+")
# ---------- Agent 等价 Streams DSL ----------
@app.agent(text_topic)
async def process(stream):
"""每收到一行文本:拆词 → 累加到 counts Table。"""
async for line in stream:
text = line.decode(errors="ignore").lower()
for word in WORD_RE.split(text):
if not word:
continue
counts[word] += 1
print(f"[count] {word} -> {counts[word]}")
# ---------- HTTP 接口:Interactive Queries 风格 ----------
@app.page("/wordcount/{word}/")
async def get_word(web, request, word: str):
"""直接查 RocksDB 拿某个词的当前计数。"""
return web.json({"word": word, "count": counts.get(word, 0)})
@app.page("/wordcount/top")
async def top_n(web, request):
"""返回 Top-10 高频词。"""
items = sorted(counts.items(), key=lambda x: -x[1])[:10]
return web.json([{"word": w, "count": c} for w, c in items])
if __name__ == "__main__":
app.main()
# ============================================================================
# 与 Java Streams 对照
# ----------------------------------------------------------------------------
# | Java Streams | Faust |
# |---------------------------------------|----------------------------------|
# | StreamsBuilder.stream("learn.text") | app.topic("learn.text") |
# | KTable + Materialized.as("xxx") | app.Table("xxx", default=int) |
# | flatMapValues + groupBy + count | for word in line.split: cnt[w]++ |
# | num.stream.threads | 启动多个 worker 进程 |
# | num.standby.replicas | table_standby_replicas |
# | EOS v2 | Faust 不支持端到端 EOS(仅 ALO) |
# | Interactive Queries | @app.page web 接口 |
# ----------------------------------------------------------------------------
# 性能对照(参考):
# - 单 worker 吞吐:Faust ≈ 5-20K msg/s(Python GIL 限制)
# - JVM Streams 单线程:50-200K msg/s
# - 高吞吐用 Streams / Flink;Python 团队的中小流量用 Faust 完全够。
# ============================================================================java
/*
* ====================================================================
* 第 17 章 - Kafka Streams 与 ksqlDB
* wordcount_streams.java - 经典 WordCount 的 Kafka Streams Java 实现
* --------------------------------------------------------------------
* Maven 依赖(pom.xml):
*
* <dependency>
* <groupId>org.apache.kafka</groupId>
* <artifactId>kafka-streams</artifactId>
* <version>3.8.0</version>
* </dependency>
*
* 准备:
* kafka-topics.sh --bootstrap-server localhost:9092 \
* --create --topic learn.text --partitions 3 --replication-factor 1
* kafka-topics.sh --bootstrap-server localhost:9092 \
* --create --topic learn.wc-output --partitions 3 --replication-factor 1
*
* 运行:
* mvn package
* java -jar target/wordcount-1.0.jar
*
* 输入:
* kafka-console-producer.sh --bootstrap-server localhost:9092 --topic learn.text
* > kafka streams kafka topic kafka
*
* 观察输出:
* kafka-console-consumer.sh --bootstrap-server localhost:9092 \
* --topic learn.wc-output --from-beginning \
* --property print.key=true \
* --value-deserializer org.apache.kafka.common.serialization.LongDeserializer
* ====================================================================
*/
package com.example.kafka.streams;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;
import java.util.Arrays;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;
public class WordCount {
public static void main(String[] args) {
// ===== 1. 配置 =====
Properties props = new Properties();
// application.id 同时是 consumer group + 内部 changelog/repartition 的前缀
// 必须全局唯一!
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 每个实例的线程数;total threads = instances × this
props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 2);
// EOS v2(可选):把消费 + 写 + State + changelog 打包成一个事务
// 代价是吞吐 ~10% 损耗、commit interval 默认 100ms 的延迟
// props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
// standby replica:另一实例热备一份 State Store,主挂秒级接管
props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);
// ===== 2. 构建 Topology =====
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> textLines = builder.stream("learn.text");
KTable<String, Long> wordCounts = textLines
// 一句话 → 多个 word(无状态算子)
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split("\\W+")))
// 过滤掉空 word
.filter((key, word) -> word != null && !word.isEmpty())
// 改 key 为 word 自己;这一步会触发后续 repartition
.groupBy((key, word) -> word)
// 聚合:count,输出是 KTable<word, count>
// Materialized.as 给 State Store 显式命名,便于 IQ 查询
.count(Materialized.as("counts-store"));
// 把 KTable 转回 KStream 并写到 sink topic
wordCounts.toStream()
.to("learn.wc-output", Produced.with(Serdes.String(), Serdes.Long()));
// ===== 3. 启动 =====
KafkaStreams streams = new KafkaStreams(builder.build(), props);
CountDownLatch latch = new CountDownLatch(1);
// 优雅关闭
Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown") {
@Override
public void run() {
streams.close();
latch.countDown();
}
});
// 状态恢复进度监听(生产建议接入)
streams.setStateListener((newState, oldState) ->
System.out.printf("[State] %s -> %s%n", oldState, newState));
// Uncaught exception 处理:默认 SHUTDOWN_CLIENT;可改 REPLACE_THREAD(线程级容错)
streams.setUncaughtExceptionHandler(ex -> {
System.err.println("Uncaught: " + ex);
return org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler
.StreamThreadExceptionResponse.REPLACE_THREAD;
});
// 打印 Topology(调试好帮手)
System.out.println(builder.build().describe());
try {
streams.start();
latch.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
System.exit(1);
}
}
}
/*
* ====================================================================
* 内部产生的 Kafka Topic(应用启动后会自动创建):
* wordcount-app-counts-store-changelog ← KTable 的 Changelog(compacted)
* wordcount-app-KSTREAM-AGGREGATE-STATE-... ← repartition topic(按 word 重分区)
*
* 想干净启动(清掉所有内部 topic + 本地 RocksDB):
* kafka-streams-application-reset.sh \
* --application-id wordcount-app \
* --input-topics learn.text \
* --bootstrap-server localhost:9092
* rm -rf /tmp/kafka-streams/wordcount-app
*
* ====================================================================
*/ksql_examples.sql ↗ · realtime_topn.py ↗ · windowed_aggregate.py ↗ · wordcount_faust.py ↗ · wordcount_streams.java ↗