像等一份分开发来的体检报告:数据、影像和说明可能先后颠倒,但三样齐了就该立刻分析,不必等下一批。原帖讨论的正是这种场景:设备会为同一时间戳分别发送 data points、frames 和 metadata。微批处理——把短时间内的数据攒成小批计算——若只按到达顺序拼接,很容易串错记录。
据 Reddit r/apachekafka 社区讨论,常见做法是让三路记录使用同一个业务键,例如“设备 ID+时间戳”,再把已到部件分别存入 state store——处理程序保存中间状态的地方。每来一件,系统就检查另外两件;配齐后立即合并、输出并删除状态。对于始终缺件的组合,可用 punctuator——定期执行的检查任务——扫描并过期;扫描越频繁,超时越准,但开销也越高。
“只处理一次”还不能靠一条完成标记口头保证。Apache Kafka 官方文档指出,Kafka Streams 默认是 at-least-once,也就是故障重试时可能重复;显式启用 exactly-once 后,才能协调输入进度、状态更新和输出写入。若采用先落入 landing/bronze 表、再周期性回看拼接的批式方案,还要与已进入 silver 表的数据去重。