Kafka Streams 把「有状态流处理」做成了库而不是框架:状态存在本地磁盘的 RocksDB 里,同时用 changelog topic 在 Kafka 上做持久化备份。这个设计的妙处在于,计算不依赖外部数据库,扩缩容与故障恢复都通过 Kafka 自身的分区机制完成;代价是状态恢复要走 changelog 重放,而恢复时间直接决定你的可用性预算。理解状态存储的机制与 RocksDB 的调优旋钮,是让 Streams 应用在生产上稳定运行的前提。
1. 状态存储的类型与选择
1.1 为什么流处理需要状态
无状态算子(map、filter、flatMap)逐条处理,不需要记住任何东西。但一旦涉及聚合、join、窗口、去重,就必须「记住过去」:
- 聚合:
count()、sum()要累加历史值。 - Join:stream-table join 要把整张表缓存在本地,才能对每条流记录做查表。
- 窗口:要按 key 维护窗口内的所有记录,并在窗口关闭后触发计算。
- 去重:要记住已经见过的 key。
这些「记住的东西」就是状态(state)。Kafka Streams 把状态显式建模为一等公民,而不是藏在算子内部的散列表,因此可以持久化、可以恢复、可以查询。
1.2 三类存储
Kafka Streams 提供三种 StateStore 抽象:
| 类型 | 接口 | 语义 | 典型用途 |
|---|---|---|---|
| KeyValueStore | KeyValueStore<K,V> | 每 key 一个值 | 聚合、join、去重 |
| WindowStore | WindowStore<K,V> | 每 key 每窗口一个值 | 时间窗口聚合 |
| SessionStore | SessionStore<K,AGG> | 每 key 一组会话 | 会话窗口、用户行为分析 |
它们的实现都基于同一个底座:RocksDB(或内存版 InMemoryKeyValueStore)。理解这一点很重要——调优的手段最终都落在 RocksDB 上。
1.3 持久化与内存的取舍
Stores.persistentKeyValueStore() 用 RocksDB 落地磁盘,Stores.inMemoryKeyValueStore() 用堆内存的 TreeMap。选择取决于状态规模:
// 持久化:状态可以远超内存,但读走磁盘
StoreBuilder<KeyValueStore<String, Long>> persistent =
Stores.keyValueStoreBuilder(
Stores.persistentKeyValueStore("counts-store"),
Serdes.String(),
Serdes.Long());
// 内存:极快,但状态必须能全部装进堆,且恢复要重放全部 changelog
StoreBuilder<KeyValueStore<String, Long>> inMemory =
Stores.keyValueStoreBuilder(
Stores.inMemoryKeyValueStore("counts-store"),
Serdes.String(),
Serdes.Long());
一个常被忽略的差异是恢复成本。内存存储不落盘,重启后必须从 changelog 的第一条开始重放;RocksDB 存储虽然也依赖 changelog,但本地已有大部分数据,只需重放「上次刷盘之后」的增量。因此生产上除非状态很小(几十 MB 以内),否则一律用持久化存储。
2. 窗口状态与 KV 状态
2.1 KeyValueStore 的语义
KV 存储是最基础的抽象。它的读写都发生在当前处理的任务(Task)内,即某个分区上:
StreamsBuilder builder = new StreamsBuilder();
KTable<String, Long> counts = builder
.stream("orders", Consumed.with(Serdes.String(), Serdes.String()))
.groupByKey()
.count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("order-counts")
.withKeySerde(Serdes.String())
.withValueSerde(Serdes.Long()));
Materialized.as("order-counts") 这一步同时做了三件事:命名状态存储、注册 changelog topic(order-counts-changelog)、把存储挂到拓扑上供交互式查询。命名不是可选的——匿名存储无法被查询,也不便于运维定位。
2.2 窗口存储与保留
窗口存储为每个 (key, windowStart) 组合存一个值。窗口的保留期(retention)决定状态何时可以被清理:
TimeWindows windows = TimeWindows
.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(30))
.advanceBy(Duration.ofMinutes(1)); // 滑动窗口的步长
KTable<Windowed<String>, Long> windowed = stream
.groupByKey()
.windowedBy(windows)
.count(Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("win-counts")
.withRetention(Duration.ofHours(2)) // 状态保留 2 小时
.withCachingEnabled());
三个时间概念必须分清:
- 窗口大小(size):窗口覆盖的时间跨度。
- 宽限期(grace):窗口结束后还允许迟到的时长。超过 grace 的记录被丢弃(
WINDOW_CLOSE语义)。 - 保留期(retention):状态在存储里保留多久。必须 ≥ 窗口大小 + 宽限期,否则会出现「窗口还没关闭,状态就被清了」的诡异丢数。
默认 retention 是 24 小时,对大窗口或高频 key 来说,状态会急剧膨胀。这是窗口作业内存/磁盘暴涨的头号原因。
2.3 去重与 session 窗口
去重 是 KV 存储的经典用法:把见过的 key 存进去,处理前先查:
KStream<String, String> deduped = stream
.transform(() -> new Transformer<String, String, KeyValue<String, String>>() {
private KeyValueStore<String, Long> seen;
@Override
public void init(ProcessorContext ctx) {
seen = ctx.getStateStore("seen-store");
}
@Override
public KeyValue<String, String> transform(String k, String v) {
if (seen.get(k) != null) return null; // 已见过,丢弃
seen.put(k, System.currentTimeMillis());
return KeyValue.pair(k, v);
}
}, "seen-store");
注意这个去重集合会无限增长,必须配 Stores.persistentKeyValueStore 的 withLoggingEnabled 与定期清理(或用 windowedBy 的窗口存储自动过期)。
SessionStore 则把「间隔小于 inactivity gap 的记录」聚成一个会话,天然适合用户行为分析。它的状态量取决于会话数与会话长度,往往比窗口存储更难预估。
3. 变更日志与故障恢复
3.1 changelog topic 机制
每个带日志的状态存储对应一个 changelog topic,命名规则是 <store-name>-changelog。它的行为特征:
- 分区数等于源 topic,保证同一 key 的状态变更落在同一分区。
- 清理策略是
compact(不是 delete),因为只需要保留每个 key 的最新值。 - 每条写入状态的操作都会同步写一条 changelog。这带来写放大:一次
count()在磁盘上是「RocksDB 写 + changelog 写」。
写放大的缓解手段是缓存(caching) 与 commit.interval.ms。开启缓存后,同一 key 的多次更新在内存里合并,只在提交时才向下游与 changelog 发一条。这能大幅减少 changelog 写入量,但会引入延迟(下游看到结果的时间变成提交间隔)。
# 缓存大小(每个线程的缓冲字节数)
cache.max.bytes.buffering=10485760
# 提交间隔:越大越省写,但结果越迟可见
commit.interval.ms=30000
3.2 恢复流程与时间估算
任务从 broker A 迁移到 broker B 时,B 上没有任何本地状态,必须恢复:
- 消费 changelog topic 对应分区的全部消息。
- 逐条写入本地 RocksDB。
- 追平到 changelog 的当前末端后,任务进入 RUNNING 并开始处理。
恢复时间是可用性的核心指标。粗略估算:
恢复时间 ≈ changelog 数据量 / 恢复吞吐
changelog 数据量 ≈ 状态大小 × 写放大系数(通常 1~3)
恢复吞吐 ≈ 受限于网络与磁盘,单任务常见 50~200 MB/s
一个 10 GB 的状态,若写放大 2 倍、恢复吞吐 100 MB/s,恢复需要约 200 秒。在这段时间里该分区不处理数据,消费滞后持续增长。这解释了为什么大状态作业的再平衡格外痛苦。
3.3 standby 副本加速恢复
num.standby.replicas 让 Kafka Streams 在其他实例上维护热备状态副本。这些副本持续消费 changelog 并更新本地 RocksDB,但不处理数据。当任务迁移过去时,本地状态已经基本就绪,只需补上很小的增量,恢复时间从分钟级降到秒级。
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-aggregator");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1); // 每个任务 1 个热备
props.put(StreamsConfig.STATE_DIR_CONFIG, "/var/lib/kafka-streams");
代价是资源翻倍:每个任务的状态被复制到另一台机器,磁盘与网络开销成倍增加。实践中按「关键作业用 1 个 standby,非关键用 0」来分配。standby 数量不能超过实例数减一,否则配置无效。
4. RocksDB 调优
4.1 关键配置项
Kafka Streams 通过 RocksDBConfigSetter 暴露 RocksDB 的配置:
public class TunedRocksDBConfig implements RocksDBConfigSetter {
@Override
public void setConfig(String storeName, Options options,
Map<String, Object> configs) {
// 每个 store 的 memtable 大小:默认 16MB,写密集可调大
options.setWriteBufferSize(64L * 1024 * 1024);
options.setMaxWriteBufferNumber(3);
// 允许同时存在的 LSM 层文件数
options.setMaxOpenFiles(-1); // -1 表示不限制
// 块缓存:读密集场景的关键
BlockBasedTableConfig tableConfig = (BlockBasedTableConfig) options.tableFormatConfig();
tableConfig.setBlockCache(new LRUCache(256L * 1024 * 1024));
tableConfig.setBlockSize(16L * 1024);
options.setTableFormatConfig(tableConfig);
// 后台压缩线程数
options.setMaxBackgroundJobs(4);
}
@Override
public void close(String storeName, Options options) {
options.close();
}
}
注册方式:
rocksdb.config.setter=com.example.TunedRocksDBConfig
4.2 内存预算怎么算
RocksDB 是 LSM-tree,写入先进内存的 memtable,写满后刷成 SST 文件,后台再合并。因此内存占用分三块:
- Write buffer(memtable):
writeBufferSize × maxWriteBufferNumber × store 数 × 线程数。这是最容易失控的一块。 - Block cache:读路径的缓存,
LRUCache大小。 - 索引与过滤器:每个 SST 文件的索引块常驻内存,
maxOpenFiles = -1时文件数越多占用越大。
一个直观的例子:50 个状态存储、每个 writeBufferSize=64MB、maxWriteBufferNumber=3,仅 memtable 就要 50 × 64 × 3 = 9.6 GB(每个 Streams 线程)。这还没算 block cache。RocksDB 调优的第一原则是先算总账,再调单个参数。
实践中的经验值:给 Streams 实例分配的内存里,堆内留出足够空间(RocksDB 用堆外内存,但 Kafka Streams 自身的缓存、反序列化在堆内),堆外按「状态大小 × 0.3」预留,剩下的交给 RocksDB 的 block cache。
4.3 压缩与写放大
RocksDB 默认用 SnappyCompression,可以换成 LZ4(更快)或 ZSTD(更省空间):
options.setCompressionType(CompressionType.LZ4_COMPRESSION);
压缩影响的是 SST 文件大小(进而影响恢复时读 changelog 之外的本地位移)与 CPU。对状态大、磁盘紧张的作业,ZSTD 能省 30%~50% 空间,代价是压缩 CPU。
写放大的另一个来源是 LSM 层的 compaction。level0FileNumCompactionTrigger(默认 4)越小,compaction 越频繁、写放大越大但读放大越小。写密集的作业可以把 level0 触发阈值调大(如 8),减少 compaction 频率。
4.4 常见 RocksDB 问题
问题一:磁盘 IO 打满,消费滞后。 多半是 compaction 与恢复同时进行。缓解:调大 maxBackgroundJobs 让 compaction 更快结束,或换更快的盘(NVMe)。
问题二:Too many open files。 maxOpenFiles = -1 时每个 SST 文件占一个 fd,状态大时轻松超过 ulimit -n。要么调大 ulimit,要么设成有限值(如 1000)让 RocksDB 自行控制。
问题三:恢复后 RocksDB 目录巨大。 删除后的 key 在 LSM 里只是墓碑标记,实际空间要等 compaction 回收。可以定期触发全量 compaction,或调小 retention 让窗口状态自然过期。
5. 再平衡下的状态迁移
5.1 任务分配与状态归属
Kafka Streams 的并行单元是 Task,一个 Task 对应一个分区。状态存储是任务私有的——order-counts 存储实际按分区拆成多个实例,每个 Task 只持有自己分区的状态。
再平衡时,协调者重新分配 Task 到实例,状态目录也要跟着走。但状态不会跨机器传输(那太慢),而是通过 changelog 重建,或用 standby 副本就近接管。因此:
- 同一个 Task 尽量回到原来的实例,可以复用本地状态,避免恢复。
- Kafka Streams 的
StickyTaskAssignor(默认)就是为此设计的,它会尽量保持任务粘性。
5.2 平滑再平衡与静态成员
传统再平衡是「全员停止、重新分配」,会造成stop-the-world。Kafka 2.4 引入的增量协作再平衡(Incremental Cooperative Rebalancing)让实例可以分批迁移,未被重新分配的任务继续处理。
Kafka Streams 2.6+ 默认使用协作式再平衡。更进一步的是静态成员(static membership):给每个实例配置固定的 group.instance.id,实例短暂重启(如滚动发布)时不会触发再平衡,因为 broker 认为它还在组内(在 session.timeout.ms 内)。
group.instance.id=streams-instance-1
session.timeout.ms=30000
这对状态大的作业意义重大——滚动重启不再触发全量恢复。
5.3 状态目录管理
state.dir 下的目录结构是:
/var/lib/kafka-streams/<application.id>/<task-id>/
├── rocksdb/<store-name>/ # RocksDB 文件
└── .checkpoint # 已刷盘的 changelog offset
.checkpoint 记录了「本地状态已经对齐到 changelog 的哪个 offset」,恢复时从这里续传。这个文件丢了就会从头重放。因此:
- 不要把
state.dir放在临时卷或容器可写层(重启即丢,等于每次全量恢复)。 - 多实例共享
state.dir时要保证实例间目录隔离。 - 容器化部署时用 PVC 或 hostPath 持久化,参考 Kubernetes 上运行 Kafka Streams 。
6. 可查询状态与交互式查询
状态存储不仅能内部使用,还能被外部查询(Interactive Queries)。查询入口是 KafkaStreams.store():
StoreQueryParameters<KeyValueStore<String, Long>> params =
StoreQueryParameters.fromNameAndType(
"order-counts", QueryableStoreTypes.keyValueStore());
// 注意:只能查本实例上存在的存储
KeyValueStore<String, Long> store = streams.store(params);
Long value = store.get("order-123");
由于状态是分区的,外部请求可能落在没有目标分区的实例上。标准做法是暴露一个 HTTP 接口,先在本地查,未命中则根据 metadataForKey 找到目标实例并转发:
KeyQueryMetadata meta = streams.queryMetadataForKey("order-counts", key, Serdes.String().serializer());
if (!meta.activeHost().equals(thisInstance)) {
// 转发到 meta.activeHost()
}
需要注意查询一致性:默认查询读的是本地 RocksDB,可能落后于 changelog。若需要「读到最新」,要么走 standby 副本(queryableStoreType 用 withPartitions),要么接受最终一致。
7. 监控指标
状态相关的关键 JMX 指标:
| 指标 | 含义 | 关注点 |
|---|---|---|
restore-consumer-records-total | 已恢复的记录数 | 恢复进度 |
restore-remaining-records | 剩余待恢复记录 | 恢复还要多久 |
rocksdb-bytes-written-total | RocksDB 写入字节 | 写放大 |
rocksdb-block-cache-hit-ratio | 块缓存命中率 | < 0.8 说明缓存偏小 |
state-store-records-total | 存储记录数 | 状态增长趋势 |
commit-total | 提交次数 | 与 commit.interval 是否一致 |
restore-remaining-records 持续不降,说明恢复卡住了(常见于 changelog 分区被限流或磁盘满)。block-cache-hit-ratio 低则说明读放大严重,该加 block cache 了。
8. 常见坑清单
- 匿名状态存储(没调
Materialized.as),无法查询也无法定位 changelog。 withRetention小于「窗口大小 + grace」,窗口未关闭状态就被清。- 去重集合不设过期,状态无限增长直到磁盘爆。
state.dir放在容器临时层,每次重启全量恢复。- RocksDB
writeBufferSize盲目调大,导致堆外内存 OOM。 maxOpenFiles = -1撞上ulimit -n,报Too many open files。- 大状态作业滚动重启未配静态成员,每次触发全量再平衡。
- 只调 standby 数量却不加磁盘,恢复快了但磁盘先满。
- 把 changelog topic 的
cleanup.policy手工改成delete,破坏恢复语义。 - 交互式查询不转发,只在本地查,命中率随分区分布骤降。
9. 总结
Kafka Streams 的状态存储是「本地 RocksDB + 远端 changelog」的组合。本地保证性能,远端保证可靠,恢复时间就是两者之间的差额。调优的抓手因此很清晰:减小状态规模(收窄 retention、及时清理)、加快恢复(standby 副本、静态成员)、控制写放大(缓存与提交间隔)、算清内存账(memtable 与 block cache 的总和)。
真正把状态管好,才能谈后面的拓扑正确性。状态恢复的窗口期也正是运维监控最该盯住的时段,相关指标与告警方法可以看 Kafka 运维监控与故障恢复 。若状态规模已超出单机承受范围,就该考虑把重状态迁移到外部存储,或改用 Apache Flink 流处理 的更大规模状态后端。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。