Learn
Kafka/15-kafka-streams

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] += 1
维度Kafka StreamsFlink
部署形态库,嵌入应用,无额外集群独立集群(或 K8s Operator)
数据源只有 KafkaKafka、CDC、文件、JDBC…任意
状态规模中小(受本地盘限制)超大状态(增量 checkpoint)
事件时间/乱序处理支持,能力一般业界最强(watermark 体系完善)
SQL无(KSQL 是另一产品)Flink SQL 成熟
团队成本Java 应用运维即可需要流计算平台运维能力

经验法则:数据进出都是 Kafka、逻辑是过滤/富化/中等规模聚合 → Streams(少一套集群就是少一堆事故);多源多汇、超大状态、复杂事件时间语义、需要 SQL 化 → Flink。

⚠️Streams 常见坑
  1. application.id 就是消费组 id 且关联所有内部 topic 名:随意改名等于换了个新应用,状态与进度全部从零开始。
  2. 状态目录 state.dir 放在容器临时盘:每次发布都全量重放 changelog,恢复极慢。K8s 上要挂持久卷或配 standby 副本。
  3. 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 是应用身份的根,改名 = 重来
🎯练习
  1. 写一个 Streams 应用:从 orders 读 JSON,过滤金额大于 100 的写入 big-orders,运行后用 kafka-topics --list 找找自动创建的内部 topic。
  2. 实现每用户 1 分钟滚动窗口的订单计数,输出到结果 topic,用 console consumer 观察窗口结果的更新方式。
  3. 起两个相同 application.id 的实例,kill 掉一个,观察任务迁移与状态恢复日志;再配 num.standby.replicas=1 对比恢复速度。