07. Kafka Streams 流处理

Kafka Streams: KStream、KTable、窗口聚合与 Exactly-Once 语义

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 提交实现。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配
  3. Kafka 详解:分布式日志系统、ISR 与一致性保证