Kafka Streams 流处理:从消费者到实时计算引擎
提出问题
"Kafka 能做流处理吗?"——这是面试里一个很能区分深浅的问题。浅的回答是"Kafka 是消息队列,流处理得用 Flink";深一点的会说"Kafka 自带 Kafka Streams,轻量场景不用额外部署 Flink 集群"。
现实中的场景很典型:你有一个订单流,需要实时统计"每分钟各城市的成交额""连续下单 3 次的用户""订单流和用户流关联出画像"。用普通 Consumer 写,你得自己维护状态(每个城市的累加值放哪?——用 ConcurrentHashMap 的话,重启就没了)、自己处理时间窗口(一分钟怎么切?——用 ScheduledExecutorService 周期归零,但怎么可能精确对齐到 00:00:00 这个整一分钟边界?)、自己扛住重启后状态不丢(进程挂了累加值怎么办?——写本地文件,但多实例选主谁来写?)。
实际踩过的坑:我用普通 Consumer + ConcurrentHashMap 做过一个"过去 5 分钟实时 PV 统计"。上线第一天就发现:Consumer 重平衡时,分区被重新分配,旧分区上的累加值全丢了;进程重启更惨,整个 5 分钟窗口数据归零,监控看板直接跳空。后来换成 Kafka Streams,这些全由 State Store + changelog 自动处理了。
面试官想确认的是:你知不知道 Kafka Streams 和普通 Consumer 的本质区别?你懂不懂有状态计算背后的 State Store 和 changelog 机制?你能不能说清它和 Flink 的取舍边界?
分析问题
一、Kafka Streams 是什么:一个库,不是一个集群
最关键的认知:Kafka Streams 是一个 Java 库(org.apache.kafka:kafka-streams),不是一个独立部署的计算集群。它就是一个 jar 包,嵌进你的 Spring Boot 应用里跑。这和 Flink(需要 JobManager + TaskManager 集群)是根本区别。
它的并行度直接靠 Kafka 的分区数撑起来:你的应用起 N 个实例,Streams 自动把 Topic 的分区分给这 N 个实例,跟消费组重平衡是同一套机制。扩容就是多起几个 Pod,不需要动任何集群配置。
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-stream-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());
// 关键配置:数十分钟内重复消费相同 offset 用不到,但开太大 state store 会膨胀
props.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 1000); // 每秒提交一次 offset,影响恢复时间
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> orders = builder.stream("orders");
// 实时统计每个城市的订单数
orders
.groupBy((key, value) -> extractCity(value))
.count()
.toStream()
.to("city-order-count", Produced.with(Serdes.String(), Serdes.Long()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();就这几行,一个实时统计的流处理任务就跑起来了——没有集群,没有 YARN,没有 JobManager。注意:APPLICATION_ID_CONFIG 不能乱改,它是 State Store 的命名空间前缀,改了之后重启会重建 RocksDB,从 changelog topic 重新回放,如果 changelog topic 保留期不够长(默认 1 天,但生产环境压缩 topic 会按 min.compaction.lag.ms 清除),恢复时间会很长甚至数据不全。
二、KStream vs KTable:流表二象性
这是 Kafka Streams 的核心抽象,也是面试深水区。
- KStream:无界的事件流,每条记录都是一个独立事件。同一个 key 出现多次,代表发生了多次(比如"用户 A 下单""用户 A 又下单"是两个事件)。相当于数据库的 append-only log。
- KTable:变更日志的物化视图,同一个 key 的新记录覆盖旧记录,代表状态的最新值(比如"用户 A 的余额",后一条覆盖前一条)。相当于数据库的 current snapshot。
// KStream:每条都是事件,累加
KStream<String, Long> clicks = builder.stream("clicks");
// KTable:同 key 覆盖,代表最新状态
KTable<String, Long> userBalance = builder.table("user-balance");
// 流表 join:给点击流补充用户余额(维表关联)
clicks.join(userBalance, (click, balance) -> enrichClick(click, balance));流表 join 的陷阱:KStream-KTable join 默认是左表驱动的——只有 KStream 侧来事件时才会触发 join 并输出。如果维表(KTable)侧的数据更新发生在 KStream 事件之后,这条 KStream 事件不会再次 join。这会导致:用户下单时余额是旧的,但订单已经记了旧余额。解决方案有两种:一是用 KTable-KTable join 做双向触发(但只适用于两边都是 KTable 的场景);二是自己实现延迟关联——订单先落库,30 秒后用 KafkaStreams#schedule 定时从外部 DB 补查维表再更新输出。
"流表二象性"是精髓:一个 KStream 聚合后变成 KTable(累加值是状态),一个 KTable 的每次变更又可以转成 KStream(每次变更是事件)。这套抽象让你能像写 SQL 一样处理流。
真实对比:我一个同事用 KStream 做 UV 去重计数,直接 stream.groupByKey().count(),结果每次分区重平衡后计数归零重算——因为 KStream 的 groupBy 默认是 no-store,不落盘。正确做法是用 KTable 的 groupByKey().count(Materialized.as("uv-store")) 指定持久化状态。
| 维度 | KStream | KTable |
|---|---|---|
| 语义 | append-only 事件流 | 更新日志的物化视图 |
| 同 key 新记录 | 保留为独立事件 | 覆盖旧记录 |
| 典型场景 | 点击流、日志流、交易流水 | 用户画像、余额、配置快照 |
| 聚合结果 | 生成为 KTable | 自身就是物化结果 |
| 持久化 | 默认不持久 | 通过 Materialized 持久化到 RocksDB |
三、有状态计算与 State Store:状态存哪、丢不丢
普通 Consumer 做累加,状态放内存里,进程一挂就全丢了。Kafka Streams 用 State Store 解决:
- 本地状态默认存在嵌入式 RocksDB(也可纯内存),聚合的中间结果实时写进去。RocksDB 是 LSM-Tree 引擎,写性能好,但读放大问题需要注意:频繁随机读聚合状态时,P99 延迟可能到 20-50ms,而纯内存实现是 <1ms。
- 关键机制:每个 State Store 背后有一个 changelog topic(Kafka 内部自动创建,名称为
{application-id}-{store-name}-changelog),状态的每次变更都同步写进这个 topic。changelog topic 是 compacted topic(按 key 压缩),只保留每个 key 的最新值,不会无限膨胀。 - 容错恢复流程:进程崩溃重启后,Streams 从 changelog topic 回放恢复 State Store,状态不丢。恢复时间大致 = changelog topic 数据量 / 单分区回放吞吐。如果你的 state store 有 50GB 数据,changelog topic 回放可能需要 5-10 分钟,这段时间内该 task 无法处理新数据。
本地聚合 → 写 RocksDB State Store → 同步写 changelog topic (Kafka compacted)
↓ 崩溃重启
从 changelog 回放恢复
(恢复时间 ≈ 数据量 / 回放速度)生产环境踩坑:我们的订单实时统计 state store 默认存 RocksDB,磁盘是普通 SSD 而非 NVMe,重启后回放 changelog 花了 12 分钟,期间该分区数据消费延迟持续飙升。后来加了两条措施:
- 用
state.dir配置指定 NVMe 盘路径(延迟降到 4 分钟) - 开启 WAL(Write-Ahead Log)同步,但牺牲一点写入吞吐
面试话术:普通 Consumer 是"无状态搬运工",Kafka Streams 是"有状态计算引擎"——差别就在这个 State Store + changelog 的容错闭环上。
State Store 的三种模式:
| 模式 | 存储方式 | 适用场景 | 注意事项 |
|---|---|---|---|
| InMemoryKeyValueStore | 内存 HashMap | 数据量小、允许重启丢失 | 不参与 changelog 容错 |
| RocksDBKeyValueStore | 本地 RocksDB | 默认模式,适合大多数场景 | 注意磁盘 IO 和恢复时间 |
| PersistentKeyValueStore | 本地文件系统 | 自定义序列化 | 性能不如 RocksDB |
四、时间窗口:怎么切一分钟
流处理绕不开时间窗口。Kafka Streams 支持四种:
- Tumbling Window(滚动窗口):固定大小不重叠,"每分钟成交额"就是它。窗口边界是固定的(00:00:00-00:01:00, 00:01:00-00:02:00...),每 60 秒一个窗口。
- Hopping Window(跳跃窗口):固定大小可重叠,"每 10 秒统计过去 1 分钟"——窗口大小 1 分钟,步长 10 秒,每个事件会落入 6 个窗口。
- Sliding Window(滑动窗口):基于事件时间差的滑动,通常用于 join 时限定关联时间范围(如"订单和支付时间差不超过 5 分钟")。
- Session Window(会话窗口):按活跃间隙切分,"用户一次会话内的行为"。如果两次事件间隔超过 session gap,就切一个新会话。
// 滚动窗口:每分钟各城市成交额
orders
.groupBy((k, v) -> extractCity(v))
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
.aggregate(
() -> 0.0,
(city, order, total) -> total + extractAmount(order),
Materialized.<String, Double, WindowStore<Bytes, byte[]>>as("city-amount-store")
.withKeySerde(Serdes.String())
.withValueSerde(Serdes.Double())
)
.toStream()
.foreach((windowedKey, amount) ->
System.out.printf("城市 %s 在窗口 [%s, %s] 成交额: %.2f%n",
windowedKey.key(),
Instant.ofEpochMilli(windowedKey.window().start()),
Instant.ofEpochMilli(windowedKey.window().end()),
amount));还要处理乱序事件:用事件时间(event-time)而非处理时间(processing-time),配合 grace period 容忍迟到数据。ofSizeWithNoGrace 表示不允许迟到数据,任何窗口关闭后的数据直接丢弃。如果业务上需要容忍 5 分钟迟到,改为:
TimeWindows.ofSizeWithGrace(Duration.ofMinutes(1), Duration.ofMinutes(5))乱序的代价:容忍 5 分钟迟到意味着窗口要在真实结束时间后 5 分钟才关闭,占用内存的时间更长。如果 QPS 是 10 万/秒,每个窗口状态 100MB,同时保留 5 分钟的迟到窗口,RocksDB 的 state store 可能膨胀到 500MB 以上。可以在 Materialized 中设置 withRetention(Duration.ofDays(1)) 来控制窗口状态的保留时间。
五、Kafka Streams 拓扑结构:DSL vs Processor API
Kafka Streams 提供两种 API 来构建处理拓扑:
- DSL API(High-Level):用
KStream、KTable、groupBy、join等声明式算子,类似 Java Stream API。上面所有例子都是 DSL API。 - Processor API(Low-Level):让你手动定义
Processor节点,通过context.forward()控制数据流向。适合需要手动控制状态、计时器、分支逻辑的复杂场景。
// Processor API 示例:手动控制事件去重
class DedupProcessor implements Processor<String, String, String, String> {
private KeyValueStore<String, Long> store;
@Override
public void init(ProcessorContext<String, String> context) {
// 获取状态存储
this.store = context.getStateStore("dedup-store");
// 注册定时器,每 10 秒清理过期 key
context.schedule(Duration.ofSeconds(10), PunctuationType.WALL_CLOCK_TIME,
timestamp -> { /* 清理过期 key */ });
}
@Override
public void process(Record<String, String> record) {
Long lastSeen = store.get(record.key());
if (lastSeen == null || System.currentTimeMillis() - lastSeen > 60000) {
store.put(record.key(), System.currentTimeMillis());
context.forward(record); // 放行
}
// 1 分钟内重复 key 直接丢弃
}
}
Topology topology = new Topology();
topology.addSource("Source", "events")
.addProcessor("Dedup", () -> new DedupProcessor(), "Source")
.addStateStore(Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("dedup-store"), Serdes.String(), Serdes.Long()), "Dedup")
.addSink("Sink", "deduped-events", "Dedup");DSL vs Processor API 的选择:70% 的场景 DSL API 就够用;需要自定义状态管理、手动定时器、多分支拓扑时用 Processor API。面试时提一句"Kafka Streams 支持 Processor API 做底层扩展,不像 Flink 的 DataStream API 那样必须走集群"能加分。
总结
Kafka Streams vs Flink 怎么选
| 维度 | Kafka Streams | Flink |
|---|---|---|
| 部署形态 | 一个库,嵌进应用 | 独立集群(JobManager/TaskManager) |
| 数据源 | 只能 Kafka(1.x 起可结合 Kafka Connect 接入其他源) | Kafka/文件/JDBC/CDC 等多源 |
| 运维成本 | 低(就是个 Spring Boot 应用) | 高(要维护集群,至少 3 台 JM + N 台 TM) |
| 状态规模 | 中小(RocksDB 本地,建议单 state store < 100GB) | 大(支持超大状态 + 增量 checkpoint,单 TM 100GB+ 常见) |
| 状态恢复时间 | 分钟级(changelog 回放) | 秒级(增量 checkpoint 的恢复时间 ≈ 最后 checkpoint 的 size / 带宽) |
| 复杂计算 | 中等(窗口、聚合、join 够用) | 强(CEP、复杂窗口、批流一体、多流 join) |
| 适用场景 | 纯 Kafka 生态、轻量实时计算 | 多源、超大状态、复杂 ETL |
| 版本兼容 | 必须与 Kafka Broker 版本匹配(建议相同 major 版本) | 独立于 Kafka 版本 |
一条决策链
数据源就是 Kafka,团队不想多维护一套集群,计算逻辑不算太重 → Kafka Streams。
多数据源、超大状态、复杂 CEP、批流一体 → Flink。
介于两者之间:先用 Kafka Streams 快速上线,当状态规模超过 100GB 或者需要多源 join 时,再逐步迁移到 Flink。Kafka Streams 和 Flink 可以共存——Kafka Streams 做轻量预处理,Flink 做重度聚合。
面试话术示例
"我们订单实时看板一开始想上 Flink,但评估下来数据源只有 Kafka,团队也没有 Flink 运维经验,就用了 Kafka Streams——它就是个 jar 包嵌在现有 Spring Boot 服务里,扩容跟着分区走,State Store 用 RocksDB + changelog 做容错,重启状态不丢。后来有个需求要关联 MySQL 的 CDC 流做多源 join,那块才迁到 Flink。选型的核心是:Kafka Streams 省运维,Flink 上限高,按数据源和状态规模划边界。
一句话总结:Kafka Streams 是给 Kafka 重度用户准备的"零额外成本"流处理方案,但它不是万能的——状态超过 100GB 或者需要多源 join 时,老老实实上 Flink。"
参考:Apache Kafka 官方文档 — Kafka Streams;《Kafka Streams in Action》;Confluent 博客 — Streams and Tables in Apache Kafka;KIP-328(Kafka Streams 2.0 改进)