Kafka Streams 入门
统计"每个用户最近 5 分钟的下单金额",用普通消费者写,你得自己管状态存储、时间窗口、故障恢复、扩缩容时的状态迁移。Kafka Streams 把这些做成了库——不是集群,就是个 Java 库,嵌进你的服务,main 方法一跑就是流处理应用。
1. 流表二元性:理解 Streams 的钥匙
同一份数据,两种视角:
流 (Stream): 记录每一次变化的事件序列 —— "发生了什么"
(alice, +100) (bob, +50) (alice, -30) (alice, +10)
表 (Table): 每个 key 的当前状态 —— "现在是什么"
alice -> 80
bob -> 50
流 -> 表: 对流按 key 不断聚合/覆盖 (aggregate)
表 -> 流: 捕获表的每次变更 (changelog)这就是流表二元性(stream-table duality):表是流的聚合快照,流是表的变更日志。Kafka 的 compact topic(第 12 章)正是"表"的物理形态。
2. 三种抽象
| 抽象 | 语义 | 数据来源 |
|---|---|---|
| KStream | 事件流,每条记录独立有意义 | 普通 topic |
| KTable | 变更日志表,同 key 后者覆盖前者,分区级(每实例只持有自己分区的部分) | compact topic |
| GlobalKTable | 每个实例持有全量数据的表 | 小数据量的维表 |
选择:订单事件用 KStream;用户余额用 KTable;几千行的城市字典要跟流做 join,用 GlobalKTable(免去重分区)。
3. 无状态与有状态算子
3.1 无状态:来一条处理一条
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> orders = builder.stream("orders");
orders
.filter((key, value) -> value.contains("\"paid\"")) // 过滤
.mapValues(v -> v.toUpperCase()) // 转换值
.selectKey((k, v) -> extractUserId(v)) // 换 key (触发重分区!)
.to("paid-orders"); // 写出3.2 有状态:需要记住历史
// 每个用户的累计下单次数 -> KTable
KTable<String, Long> orderCounts = orders
.groupByKey()
.count(Materialized.as("order-counts-store")); // 命名状态存储, 可查询
// 每用户 5 分钟滚动窗口的下单金额
KTable<Windowed<String>, Double> amounts5m = orders
.mapValues(v -> parseAmount(v))
.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5)))
.aggregate(() -> 0.0, (k, amt, agg) -> agg + amt,
Materialized.as("amount-5m-store"));
// 流表 join: 给订单流补上用户等级
KTable<String, String> userLevels = builder.table("user-levels");
orders.join(userLevels,
(orderJson, level) -> orderJson + ",\"level\":\"" + level + "\"")
.to("orders-enriched");窗口类型速查:
| 窗口 | 特点 | 例子 |
|---|---|---|
| Tumbling 滚动 | 固定大小、不重叠 | 每 5 分钟一个统计段 |
| Hopping 跳跃 | 固定大小、按步长重叠 | 每 1 分钟输出最近 5 分钟 |
| Sliding 滑动 | 由相邻事件间隔定义 | 两事件相距 10s 内算同窗 |
| Session 会话 | 按不活跃间隙切分 | 用户一次访问会话 |
3.3 启动骨架
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-analytics"); // 即消费组 id
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2); // 精确一次
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));扩容 = 多起几个相同 application.id 的实例,任务(按输入分区切分)自动重分配——复用的正是消费组机制。
4. 状态存储与容错
有状态算子的状态存在本地 RocksDB,同时异步备份到 Kafka 的 changelog topic(compact):
实例 A Kafka
RocksDB [order-counts] --备份--> order-analytics-order-counts-changelog
| |
实例 A 宕机 |
v v
实例 B 接管该分区任务 <--重放恢复状态----------+- 恢复大状态可能要重放很久,
num.standby.replicas=1让另一实例维护热备状态副本,接管秒级完成。 - EXACTLY_ONCE_V2 用事务(第 6 章)把"消费 offset + 状态 changelog + 输出消息"原子提交,实现端到端精确一次。
selectKey/groupBy 改 key 后需要重分区(repartition):Streams 自动创建内部 repartition topic 把同 key 数据汇到同一任务——这是 Streams 应用会"凭空多出一些 topic"的原因,属正常现象。
5. Python 生态说明
Kafka Streams 是纯 Java/Scala 库,Python 没有官方对应。Python 技术栈的选择:
- Faust(社区维护的 faust-streaming):API 风格类似 Streams,适合轻量场景;
- 直接用 confluent-kafka 手写消费-处理-生产循环 + 外部状态(Redis);
- 重活交给 Flink(PyFlink 支持较完善)。
# faust 风格示意: 每用户订单计数
import faust
app = faust.App("order-analytics", broker="kafka://localhost:9092")
orders_topic = app.topic("orders", key_type=str, value_type=str)
counts = app.Table("order_counts", default=int)
@app.agent(orders_topic)
async def count_orders(stream):
async for key, _value in stream.items():
counts[key] += 16. Streams 还是 Flink?
| 维度 | Kafka Streams | Flink |
|---|---|---|
| 部署形态 | 库,嵌入应用,无额外集群 | 独立集群(或 K8s Operator) |
| 数据源 | 只有 Kafka | Kafka、CDC、文件、JDBC…任意 |
| 状态规模 | 中小(受本地盘限制) | 超大状态(增量 checkpoint) |
| 事件时间/乱序处理 | 支持,能力一般 | 业界最强(watermark 体系完善) |
| SQL | 无(KSQL 是另一产品) | Flink SQL 成熟 |
| 团队成本 | Java 应用运维即可 | 需要流计算平台运维能力 |
经验法则:数据进出都是 Kafka、逻辑是过滤/富化/中等规模聚合 → Streams(少一套集群就是少一堆事故);多源多汇、超大状态、复杂事件时间语义、需要 SQL 化 → Flink。
⚠️Streams 常见坑
- application.id 就是消费组 id 且关联所有内部 topic 名:随意改名等于换了个新应用,状态与进度全部从零开始。
- 状态目录 state.dir 放在容器临时盘:每次发布都全量重放 changelog,恢复极慢。K8s 上要挂持久卷或配 standby 副本。
- join 两个 topic 分区数不一致:KStream-KStream join 要求 co-partitioned(同分区数同 key 策略),不满足会报错或需要手动重分区。
小结
- 流表二元性:表是流的快照,流是表的 changelog;KStream/KTable/GlobalKTable 三种抽象
- 无状态算子直来直去;有状态算子(count/aggregate/join/窗口)依赖本地 RocksDB + changelog 备份
- EXACTLY_ONCE_V2 基于事务实现端到端精确一次;扩容复用消费组机制
- Streams 是库不是集群,Kafka 进出的中等复杂度处理首选;多源大状态选 Flink
- application.id 是应用身份的根,改名 = 重来
🎯练习
- 写一个 Streams 应用:从 orders 读 JSON,过滤金额大于 100 的写入 big-orders,运行后用 kafka-topics --list 找找自动创建的内部 topic。
- 实现每用户 1 分钟滚动窗口的订单计数,输出到结果 topic,用 console consumer 观察窗口结果的更新方式。
- 起两个相同 application.id 的实例,kill 掉一个,观察任务迁移与状态恢复日志;再配 num.standby.replicas=1 对比恢复速度。