引言
做实时计算,大家第一反应往往是 flink。但如果你的场景没那么重,不想维护庞大的 flink 集群,kafka streams 绝对是个被低估的利器。它直接嵌在 java 进程里,跟 kafka 亲儿子一样,部署起来就是个普通的 spring boot 应用。
最近带团队搞了几个流处理项目,把 kafka streams 从基础用法到生产级踩坑都趟了一遍。今天就把状态存储、窗口计算、exactly-once 这些核心玩法,还有多租户和运维的实战经验盘一盘。
1. 核心概念与拓扑设计:别被名词唬住
kafka streams 的核心抽象就两个:kstream 和 ktable。
说白了,kstream 就是流水账,来一条处理一条;ktable 是账本,只记最新状态(类似数据库的表)。
流处理拓扑(topology) 就是计算逻辑,本质上是个有向无环图。设计拓扑时,source 进数据,processor 搞计算,sink 吐结果。平时写 dsl 链式调用(stream().filter().map().to())挺爽,但逻辑一复杂,还是得老老实实切回 processor api,不然调试起来能让人怀疑人生。
2. spring boot 集成:自动配置虽好,参数得自己捏
spring boot 把 kafka streams 封装得很省事,加个 @enablekafkastreams 就能跑。但自动配置归自动配置,有些生产环境的参数必须自己配置,不能全指望默认值。
@configuration
@enablekafkastreams
public class kafkastreamsconfig {
@bean
public streamsbuilderfactorybeancustomizer streamscustomizer() {
return factory -> {
// 自定义配置,顺便把未捕获异常处理器也配了,防止线程默默死掉
factory.setstreamsconfiguration(customstreamsconfig());
factory.setuncaughtexceptionhandler((thread, exception) -> {
log.error("stream thread {} 崩了", thread.getname(), exception);
// 这里可以触发告警,或者直接调 kafkastreams.close() 让应用重启
});
};
}
@bean
public kafkastreamsconfiguration customstreamsconfig() {
map<string, object> props = new hashmap<>();
props.put(streamsconfig.application_id_config, "order-stream-app");
props.put(streamsconfig.bootstrap_servers_config, "localhost:9092");
// 生产环境无脑上 eos v2,老版本的 exactly_once 会让 broker 连接数爆炸
props.put(streamsconfig.processing_guarantee_config, streamsconfig.exactly_once_v2);
// 状态目录别放在系统盘,挂载独立数据盘
props.put(streamsconfig.state_dir_config, "/data/kafka-streams");
// 注意:这个异常处理器是 spring-kafka 提供的,需要确保引入了对应依赖
props.put(streamsconfig.default_deserialization_exception_handler_class_config,
sendtodeadlettertopicexceptionhandler.class);
return new kafkastreamsconfiguration(props);
}
}
spring boot 会自动管理 streamsbuilder 和 kafkastreams 的生命周期。应用启动时初始化拓扑,关闭时优雅提交 offset 并刷盘。
3. 状态存储 rocksdb:快是真快,爆也是真爆
kafka streams 敢做有状态计算,底气就在 rocksdb。数据存本地磁盘,读写极快。但稍微不注意,就能把机器内存和磁盘撑爆。
状态恢复与 standby replicas
状态存储不是孤立的,kafka streams 会在后台为每个状态存储建一个内部的 changelog topic。数据改了,同步写 changelog。节点宕机重启,从 changelog 拉数据重建本地状态。
这里有个实战经验:一定要配 standby replicas(备用副本)。
props.put(streamsconfig.num_standby_replicas_config, 1);
不然节点一挂,新节点拉起时从 changelog 狂拉数据,那恢复时间能让你等到怀疑人生。配了备用副本,其他节点会异步维护一份只读副本,主节点挂了直接热切换,恢复时间从分钟级降到毫秒级。
4. 窗口聚合:关掉 grace period 保平安
流数据无限长,必须切块算。kafka 给了三种窗口:滚动、滑动、会话。
// 1. 滚动窗口:固定大小不重叠,比如每小时统计一次
ktable<windowed<string>, long> hourlyclicks = clickstream
.groupbykey()
.windowedby(timewindows.ofsizewithnograce(duration.ofhours(1)))
.count(materialized.as("hourly-clicks-store"));
// 2. 滑动窗口:固定大小但重叠,算移动平均
ktable<windowed<string>, long> movingavg = eventstream
.groupbykey()
.windowedby(timewindows.ofsizewithnograce(duration.ofminutes(5)).advanceby(duration.ofminutes(1)))
.count();
// 3. 会话窗口:动态大小,基于活动间隔,算用户在线时长
ktable<windowed<string>, long> sessionduration = useractivitystream
.groupbykey()
.windowedby(sessionwindows.ofinactivitygapwithnograce(duration.ofminutes(30)))
.count();
重点提一句 withnograce。老版本默认有 24 小时的 grace period,用来处理迟到数据。但在很多业务里,这 24 小时的宽限期会导致状态存储里堆积大量过期数据,直接把 rocksdb 撑爆。只要业务能容忍极少量的迟到数据丢失,果断加上 withnograce 关掉它。
5. 表流 join:注意单向触发和序列化大坑
把 kstream 和 ktable 做 join 是常态,比如订单流关联用户信息表。
// 1. 构建用户信息 ktable
ktable<string, userdetail> usertable = builder.table(
"user-details-topic",
consumed.with(serdes.string(), userdetailserde)
);
// 2. 订单流 join 用户表
kstream<string, enrichedorder> enrichedorders = orderstream
.leftjoin(
usertable,
(order, user) -> new enrichedorder(order, user),
// 坑点:这里的 serde 别瞎写 null,左右两边的 value serde 都得老老实实传进去
joined.with(serdes.string(), orderserde, userdetailserde)
);
两个大坑:
- 流表 join 是单向触发的。只有 kstream 来了新数据,才会去 ktable 里查;ktable 数据更新了,不会主动去推 kstream。
joined.with里的 serde 必须传全,之前见过有人右侧 value serde 传null,运行时序列化直接报错。
6. exactly-once 语义:事务打包,拒绝重复
eos (exactly-once) 听起来很玄乎,其实底层就是靠事务。
以前用 exactly_once,每个 task 一个事务生产者,broker 连接数直接爆炸。现在无脑上 exactly_once_v2,所有 task 共享一个事务生产者,资源消耗断崖式下降。
底层逻辑很简单:读数据、改状态、写结果、提交 offset,这四步打包成一个 kafka 事务。要么全成功,要么全回滚。宕机重启?没关系,从上次提交的 offset 接着跑,状态靠 changelog 恢复,数据绝对不会重复处理。
7. 交互式查询 (iq):缓存带来的一致性陷阱
iq 是个好东西,能把 streams 应用变成分布式 kv 数据库,直接通过 rest 接口查状态。
@restcontroller
@requestmapping("/api/state")
public class statequerycontroller {
// 注意:注入的是 kafkastreams,不是 streamsbuilder
private final kafkastreams kafkastreams;
public statequerycontroller(kafkastreams kafkastreams) {
this.kafkastreams = kafkastreams;
}
@getmapping("/user/{userid}/score")
public responseentity<long> getuserscore(@pathvariable string userid) {
readonlykeyvaluestore<string, long> store = kafkastreams.store(
storequeryparameters.fromnameandtype("user-score-store", queryablestoretypes.keyvaluestore())
);
long score = store.get(userid);
return score != null ? responseentity.ok(score) : responseentity.notfound().build();
}
}
巨坑预警:kafka streams 默认开了内存缓存!你刚写进去的数据,可能还在缓存里没刷到 rocksdb,这时候去查,根本查不到。
怎么解?
- 关掉缓存(性能掉底,不推荐)。
- 业务上接受最终一致性(推荐,容忍几秒延迟)。
- 如果是分布式部署,查不到数据还得自己写逻辑,通过
streamsmetadata把请求路由到真正持有那个 key 的节点上去。
8. 拓扑优化:dto 不可变与自定义序列化
dto 尽量用 lombok 的 @value 搞成不可变对象,流处理里最怕状态被意外篡改。
@value
@builder
public class orderaggregate {
string orderid;
bigdecimal totalamount;
int itemcount;
}
dsl 搞不定的复杂逻辑(比如定时任务、多状态存储联动),别硬憋,直接上 processor api。
序列化方面,json 开发快,但性能和体积不如 protobuf。生产环境数据量大的话,老老实实切 protobuf,能省下不少带宽和磁盘。
9. 多租户隔离:防住“吵闹的邻居”
saas 场景下搞多租户,最简单的就是 topic 加前缀(如 tenanta.orders)。但如果某个租户数据量特别大,把资源吃光了,其他租户就得跟着遭殃。
这时候就得物理隔离,给大租户单独分配 application id:
props.put(streamsconfig.application_id_config, "stream-app-" + tenantid);
这样每个租户拥有独立的 consumer group 和状态存储。资源配额方面,streams 自己不管这事,得靠 k8s 的 limit 和 request 来卡脖子,限制每个 pod 的 cpu 和内存。
10. 生产运维:重置、监控与异常兜底
应用重置
跑飞了或者要重跑历史数据,用重置脚本。记住,跑这脚本前必须先把应用停了!
kafka-streams-application-reset.sh --application-id order-stream-app \ --bootstrap-server localhost:9092 --input-topics orders --to-earliest
监控指标
streams 自带的 jmx 指标很全,接个 micrometer 打到 prometheus 就行。重点盯 task-closed-rate,如果这指标忽高忽低,说明 rebalance 风暴来了,赶紧去查是不是有节点在频繁挂掉或者处理太慢。
异常兜底
生产环境脏数据是常态,必须做兜底:
// 1. 反序列化遇到脏数据,扔到死信队列,别让应用直接崩了
props.put(streamsconfig.default_deserialization_exception_handler_class_config,
sendtodeadlettertopicexceptionhandler.class);
// 2. 生产异常(如 broker 拒绝写入),配置继续运行(kafka 2.8+ 支持)
props.put(streamsconfig.default_production_exception_handler_class_config,
continueonproductionexceptionhandler.class);
写在最后
kafka streams 在 java 生态里是个很实在的流处理框架,没有 flink 那么重,但能解决 80% 的实时计算需求。
当然,它也有自己的脾气,比如 rebalance 慢、状态存储调优麻烦、iq 查询有缓存延迟。用的时候得多留个心眼,别把它当成无所不能的银弹。
以上就是springboot集成apache kafka streams构建流处理应用的代码详解的详细内容,更多关于springboot apache kafka streams流处理应用的资料请关注代码网其它相关文章!
发表评论