1. 核心抽象
KStream: 每个事件独立 (插入流)
key=order1, value=created
key=order1, value=paid
key=order1, value=shipped
→ 三条独立记录
KTable: 按 key 更新 (更新流)
key=user1, value=balance=100
key=user1, value=balance=80 (更新)
→ 只保留最新状态 (log-compacted topic)
2. 窗口聚合
stream.groupByKey()
.windowedBy(TimeWindows.of(Duration.ofMinutes(5)).grace(Duration.ofMinutes(1)))
.aggregate(0L, (k, v, agg) -> agg + v.getAmount(),
Materialized.with(Serdes.String(), Serdes.Long()))
.toStream()
.to("hourly-sales");
窗口类型:
- Tumbling: 固定间隔,不重叠
- Hopping: 固定间隔,可重叠
- Session: 动态,由活动间隙定义
3. Exactly-Once
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);
通过事务 Producer + 原子 offset 提交实现。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。