Skip to content

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,不需要动任何集群配置。

java
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
java
// 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,不落盘。正确做法是用 KTablegroupByKey().count(Materialized.as("uv-store")) 指定持久化状态。

维度KStreamKTable
语义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 分钟,期间该分区数据消费延迟持续飙升。后来加了两条措施:

  1. state.dir 配置指定 NVMe 盘路径(延迟降到 4 分钟)
  2. 开启 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,就切一个新会话。
java
// 滚动窗口:每分钟各城市成交额
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 分钟迟到,改为:

java
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):用 KStreamKTablegroupByjoin 等声明式算子,类似 Java Stream API。上面所有例子都是 DSL API。
  • Processor API(Low-Level):让你手动定义 Processor 节点,通过 context.forward() 控制数据流向。适合需要手动控制状态、计时器、分支逻辑的复杂场景。
java
// 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 StreamsFlink
部署形态一个库,嵌进应用独立集群(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 改进)

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。