Skip to content

第 17 章 Kafka Streams 与 ksqlDB:让 Kafka 自己处理流

目标读者:每次要做「实时统计 / Top-N / 实时大屏 / 实时风控」就拉一个 Flink 集群、写一堆 YAML,被 Flink 的部署和状态后端折腾得崩溃;或者干脆在 Consumer 里自己 hash table 聚合,结果重启后状态全丢的同学。

学完你会:能讲清「流-表二元性」、能在 5 分钟内写出 WordCount + 1 分钟滚动窗口 GMV 统计;能区分 KStream / KTable / GlobalKTable;能解释 Streams 怎么用 RocksDB + Changelog Topic 让状态既快又可恢复;能用 ksqlDB 写 SQL 跑流处理;能讲清楚 Streams 与 Flink 的本质区别。


0. 一个生活类比:电影 vs 截图

理解 Kafka Streams 之前,记住一个最重要的隐喻:

  • 流(KStream) = 一部正在播放的电影(连续的事件序列:「每一帧画面」);
  • 表(KTable) = 这部电影任意一刻按下暂停键的截图(按 key 聚合后的当前状态)。

电影完整记录了「每一帧」(不可改变的事实),截图就是某一帧的瞬时画面。两者可以无损转换

  • 电影 → 截图:从头到尾「重放」每一帧,最后一帧就是截图(聚合 / 物化);
  • 截图 → 电影:把每张截图按时间排起来连播(Changelog)。

Kafka Streams 的核心创新就是:让流 (Topic) 和表 (State Store) 在你的代码里像同一种东西一样无缝转换。 这就是「流-表二元性」(Stream-Table Duality)。


Kafka Streams 是一个 Java 库(不是集群、不是服务、不是框架),把它 import 进你自己的 Java 应用,应用启动 = 流处理任务启动;应用关停 = 任务停。无需 Flink JobManager / Spark Driver 这种「外部调度器」。

1.1 与其他流处理引擎对比

维度Kafka StreamsksqlDBApache FlinkSpark Structured Streaming
形态Java 库SQL 引擎 + REST集群集群
部署你自己 jar 部署docker / k8sJobManager + TaskManagerDriver + Executors
调度无外部调度自己集群自带自带
状态RocksDB + Changelog同 StreamsRocksDB + DFS SnapshotRocksDB
吞吐★★★★★★★★★★★★★★★★
延迟毫秒毫秒毫秒秒(micro-batch)
上手Java DSL,门槛低SQL,最低复杂复杂
水平扩展加实例(自动 Rebalance)加实例加 TaskManager加 Executor
与 Kafka 耦合极强(仅 Kafka)极强(仅 Kafka)弱(多源)弱(多源)

核心定位:Streams = 「最轻量级的流处理」,把流处理能力嵌入到普通微服务里;Flink 是「最强大的流处理集群」,独立运维。

一句话选型:

  • 只用 Kafka做事件源 + 想跑流处理 + 不想多维护一套集群 → Kafka Streams / ksqlDB
  • 你需要接 MySQL / Pulsar / 文件等多种 source、需要 CEP / 复杂窗口 / 跨源 join → Flink

1.2 一段最简单的 Kafka Streams 代码

java
StreamsBuilder builder = new StreamsBuilder();
builder.<String, String>stream("learn.input")
       .mapValues(v -> v.toUpperCase())
       .to("learn.output");
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

这就是一个完整的流处理应用:从 learn.input Topic 读 → 转大写 → 写到 learn.output没有 JobManager、没有 Slot、没有 YAML。打成 jar 部署到 Kubernetes / Docker,加副本就是「水平扩展」。


2. 流-表二元性:Stream-Table Duality

这是整个 Streams 的灵魂。完整图示:

更精确的说法:

对象数据语义物理存储等价对象
KStream事件流:每条记录都是「发生了一件事」Kafka Topic数据库的 binlog
KTable当前状态:每个 key 的最新值本地 RocksDB + Kafka Compacted Topic(Changelog)数据库的「当前表」
GlobalKTable当前状态:每个 key 的最新值,全量复制到每个实例同 KTable,但每个 instance 全量维度表的「全表广播」

流和表的转换

text
KStream.toTable()            // 流 → 表(按 key 取最新值)
KTable.toStream()            // 表 → 流(每次更新发一条变更)
KStream.groupByKey().count() // 流 → 聚合 → 表

2.1 KStream / KTable / GlobalKTable 区别

维度KStreamKTableGlobalKTable
语义事件序列按 key 的当前值按 key 的当前值,全量
重复 key各自独立事件后值覆盖前值同 KTable
分区按 key 分区按 key 分区每实例全量副本
容量仅看 Topic 大小单实例内存 + RocksDB必须全部塞进单实例
用途业务事件、点击流累加、维度的最新状态小维度表(万级行)
Join 限制与流 Join 必须共分区与流 Join 必须共分区不要求共分区(万能 Join 利器)

记忆口诀:「事件用 KStream、状态用 KTable、广播维度表用 GlobalKTable」。


3. 拓扑(Topology):Streams 的「电路图」

Streams 程序在内部会被编译成一个 DAG(有向无环图),叫 Topology。Source 节点 → 算子节点 → Sink 节点。

可以用 streams.toString() 或调试模式打出来:

text
Topologies:
   Sub-topology: 0
    Source: source [topics: [learn.text]]
      --> flatmap
    Processor: flatmap (stores: [])
      --> groupby
    Processor: groupby (stores: [])
      --> count
    Processor: count (stores: [counts-store])
      --> sink
    Sink: sink (topic: learn.wc-output)

3.1 Task 与 Partition

  • 一个 Topology 可以拆成多个 Sub-Topology(用 repartitiongroupBy 切分);
  • 每个 Sub-Topology × 每个分区 = 一个 Stream Task
  • Task 是 Streams 的最小调度单元,一个 Task 由一个线程独占
  • 多个实例(Streams app 副本)通过 Consumer Group 协议 Rebalance 分配 Task。
text
Topic learn.text  partitions=6   ⇒  6 Stream Task
启 3 个实例,每个 num.stream.threads=2 ⇒ 6 个线程刚好一人一个 task
启 6 个实例,每个 num.stream.threads=1 ⇒ 一样 6 个 task,每实例 1 个

水平扩展只需「加实例」,Rebalance 后任务自动迁移;状态 Store 也会通过 Changelog 在新实例上重建。


4. 无状态算子:流的基本工具箱

无状态 = 处理一条消息只看这条消息本身,不依赖任何累计。

算子说明示例
map同时改 key 和 value(k,v) -> KeyValue.pair(k.toUpperCase(), v)
mapValues只改 value(不破坏分区)v -> v.toUpperCase()
filter留下满足条件的(k,v) -> v.amount > 100
filterNot反向 filter(k,v) -> v.status.equals("CANCELLED")
flatMap一变多拆 sentence → 词
flatMapValues一变多(仅 value)s -> Arrays.asList(s.split(" "))
branch / split按条件分流到多个流按金额分大单 / 小单
peek副作用(日志、metrics)不改流(k,v) -> log.info(...)
foreach终止操作,仅副作用
merge把两个 KStream 合成一个
selectKey改 key(会触发 repartition)(k,v) -> v.userId

重要规则改了 key 的算子(map / selectKey)会让流被打上 repartition 标记,下次有状态算子(groupBy / join)执行时,Streams 会偷偷在内部插一个 repartition Topic,把数据重新按新 key Hash 分区。这就是为什么有些应用启动后会发现 Kafka 多了几个 *-repartition Topic。


5. 有状态算子:流处理的「魂」

有状态 = 处理一条消息时要读 / 写一个累计的「状态」。背后必有 State Store + Changelog Topic 支撑。

5.1 groupByKey vs groupBy

java
// 已经按 key 分区,无需 repartition
KGroupedStream<String, Order> g1 = orders.groupByKey();

// 改了 key,会插 repartition
KGroupedStream<String, Order> g2 = orders.groupBy((k, v) -> v.userId);

5.2 count / reduce / aggregate

最常见的三个聚合:

java
KGroupedStream<String, Order> grouped = orders.groupByKey();

// 1) count:每 key 计数 → KTable<String, Long>
KTable<String, Long> cnt = grouped.count();

// 2) reduce:把同 key 的 value 累计成一个相同类型的值
KTable<String, Order> latest = grouped.reduce((acc, cur) -> cur);  // 取最新

// 3) aggregate:可以改输出类型;最强大、最常用
KTable<String, Stats> stats = grouped.aggregate(
    () -> new Stats(0, 0.0),                          // initializer
    (key, order, agg) -> agg.update(order.amount),    // adder
    Materialized.<String, Stats, KeyValueStore<Bytes,byte[]>>as("stats-store")
                .withValueSerde(StatsSerde)
);

5.3 join

Join 类型必须 window必须共分区典型用例
KStream-KStream✅ 必须订单流 ⨝ 支付流(10 分钟内匹配)
KStream-KTable订单流 ⨝ 用户表(流来一条查一次)
KStream-GlobalKTable流 ⨝ 国家维度表(小表广播)
KTable-KTable用户表 ⨝ 等级表(双表都按 key 维护)

KStream-KStream join 示例

java
KStream<String, Order> orders = ...;
KStream<String, Payment> payments = ...;

KStream<String, OrderWithPayment> joined = orders.join(
    payments,
    (o, p) -> new OrderWithPayment(o, p),
    JoinWindows.of(Duration.ofMinutes(10))
);

KStream-KTable join 示例

java
KTable<String, User> users = builder.table("learn.users");
KStream<String, OrderEnriched> enriched = orders.join(
    users,
    (order, user) -> new OrderEnriched(order, user)
);

5.4 窗口(Window):「用时间切流」

Streams 提供 4 种窗口:

窗口类型时长是否重叠是否对齐典型用例
Tumbling(滚动)固定❌ 不重叠✅ 时间对齐每分钟 GMV
Hopping(跳跃 / 滑动)固定 + step✅ 重叠✅ 对齐每 10s 输出最近 1 分钟
Sliding(事件级滑动)固定按事件反欺诈(最近 N 秒交易笔数)
Session(会话)不定用户活跃 session(30min 无操作切断)

生活类比:考勤表

  • Tumbling = 「每天 9:00-18:00 的考勤」—— 严格不重叠;
  • Hopping = 「每小时统计最近 4 小时考勤」—— 重叠;
  • Session = 「员工进门到离门算一次会话」—— 自动伸缩;
  • Sliding = 「每次有人刷卡就算最近 30 分钟的人数」—— 事件触发。

代码示例:

java
// 1 分钟 Tumbling
KTable<Windowed<String>, Long> perMinuteCnt = orders
    .groupByKey()
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)))
    .count();

// 1 分钟 Hopping,每 10 秒更新一次
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(5))
    .advanceBy(Duration.ofSeconds(10)))

// 30 分钟 Session
.windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(30)))

6. State Store + Changelog:状态可恢复的核心

每个有状态算子背后都有一个 State Store(默认 RocksDB),存的是「当前 key 的值」。

text
内存(RocksDB MemTable / BlockCache)

本地磁盘(RocksDB SST 文件)

Kafka Compacted Topic(Changelog,KAFKA-STREAMS-{app-id}-{store}-changelog)

6.1 写路径

每次 State Store 更新(put / delete),Streams 会同步写一条 changelog Topic(compaction 模式),保证「磁盘 + Kafka」两份。

6.2 读路径

正常读:直接从 RocksDB 读,毫秒级。

6.3 故障恢复

实例宕机 / 迁移 → 新实例从 changelog Topic 重放,把 State Store 重建。这就是「状态可恢复」的本质:状态本身就是一个 compacted topic。

6.4 standby replicas

为加速故障恢复,可以配 num.standby.replicas=1:每个 Task 在另一个实例上保持一份热备的 State Store,正常时从 changelog 持续追,主实例挂掉后秒级接管。

properties
num.standby.replicas=1

📌 与 Flink 对比:Flink 也用 RocksDB,但状态的 checkpoint 周期写到外部 DFS(HDFS / S3);Streams 没有 DFS 依赖,状态完全靠 Kafka Topic 持久化,运维更轻。


7. 时间语义:Event Time / Processing Time / Ingestion Time

时间取自适合
Event Time消息 payload 里的业务时间戳(如 order.createdAt业务真实时间,首选
Processing TimeStreams 处理消息时的系统时间简单的实时统计、不在意延迟
Ingestion TimeProducer 把消息写进 Kafka 时由 broker 盖的时间戳业务没有时间字段时的备选

7.1 自定义 TimestampExtractor

java
public class OrderEventTimeExtractor implements TimestampExtractor {
    @Override
    public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
        Order o = (Order) record.value();
        return o == null ? partitionTime : o.getCreatedAt();
    }
}

// 注册
props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG,
          OrderEventTimeExtractor.class);

7.2 乱序与水位线

业务时间是不可避免会乱序的。Streams 用 stream time(截至此刻所有消息时间戳的最大值)当水位线,超过窗口 + grace period 的迟到消息会被丢弃。

java
TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofMinutes(5))

Flink 用 Watermark,Streams 用 stream time + grace period,思路一致但 API 不同。


8. EOS v2:Streams 中的 Exactly Once

把第 13 章的事务 Producer + 幂等 Producer + read_committed 全部封装成一个开关:

java
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
          StreamsConfig.EXACTLY_ONCE_V2);

启用后:

  • 消费 source topic + 写 sink topic + 更新 State Store + 写 changelog 打包成一个事务
  • 提交一次性原子完成;
  • 中途崩溃 → 整个事务 abort,下次重做;
  • 下游消费者必须 isolation.level=read_committed 才能看不到 abort 的中间状态。

EOS v2(Streams 2.6+)相比 v1,事务 Producer 数从「每 Task 1 个」减少到「每实例 1 个」,吞吐损耗从 ~30% 降到 ~10%。

代价:吞吐 / 延迟略降;commit interval(默认 100ms)越小延迟越低但开销越大。


9. 完整示例:WordCount

经典的 Hello World:

java
public class WordCount {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();
        builder.<String, String>stream("learn.text")
               .flatMapValues(v -> Arrays.asList(v.toLowerCase().split("\\W+")))
               .filter((k, v) -> v != null && !v.isEmpty())
               .groupBy((k, v) -> v)
               .count(Materialized.as("counts-store"))
               .toStream()
               .to("learn.wc-output", Produced.with(Serdes.String(), Serdes.Long()));

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
        streams.start();
    }
}

完整版(含 Maven 依赖、Serdes、运行说明)见 17_streams/code/wordcount_streams.java

9.1 Python 等价方案:Faust

Kafka Streams 是 JVM 库,Python 没有官方等价。社区有 faust-streaming(前 Robinhood faust 的 fork),用 asyncio 风格提供类似 API:

python
import faust

app = faust.App('wordcount-py', broker='kafka://localhost:9092')
text_topic = app.topic('learn.text', value_type=str)
counts = app.Table('counts-store', default=int)

@app.agent(text_topic)
async def process(stream):
    async for line in stream:
        for word in line.lower().split():
            counts[word] += 1
            print(f'{word} -> {counts[word]}')

Faust 不与 Streams 协议互通,状态后端是 RocksDB(基于 python-rocksdb),Changelog 也叫 <app-id>-<table>-changelog 但格式独立。生产环境 Faust 适合小规模 / Python 团队,不推荐替代 JVM Streams 跑高吞吐核心链路。


10. ksqlDB:用 SQL 写流处理

ksqlDB = Kafka Streams 的 SQL 外壳:你写 SQL,ksqlDB 把 SQL 编译成 Streams 拓扑,在 ksqlDB 集群里跑。

10.1 安装与启动(一行 Docker)

bash
docker run -d -p 8088:8088 \
  -e KSQL_BOOTSTRAP_SERVERS=broker:9092 \
  -e KSQL_LISTENERS=http://0.0.0.0:8088 \
  -e KSQL_KSQL_SERVICE_ID=ksql-service \
  confluentinc/ksqldb-server:0.29.0

10.2 创建 STREAM 和 TABLE

sql
-- 把 Kafka Topic 当 STREAM 暴露(KSTREAM)
CREATE STREAM clicks (
    user_id VARCHAR KEY,
    page    VARCHAR,
    ts      BIGINT
) WITH (
    KAFKA_TOPIC='learn.clicks',
    VALUE_FORMAT='JSON',
    TIMESTAMP='ts'
);

-- 把 Kafka Topic 当 TABLE(KTABLE,按 key 取最新值)
CREATE TABLE users (
    user_id VARCHAR PRIMARY KEY,
    name    VARCHAR,
    country VARCHAR
) WITH (
    KAFKA_TOPIC='learn.users',
    VALUE_FORMAT='JSON'
);

10.3 流处理 SQL

sql
-- 1) 简单过滤 + 投影(持续查询)
CREATE STREAM tech_clicks AS
SELECT user_id, page
FROM clicks
WHERE page LIKE '/tech/%'
EMIT CHANGES;

-- 2) 1 分钟 Tumbling 窗口,按用户聚合
CREATE TABLE per_user_per_min AS
SELECT
    user_id,
    COUNT(*) AS click_cnt
FROM clicks
WINDOW TUMBLING (SIZE 1 MINUTE)
GROUP BY user_id
EMIT CHANGES;

-- 3) Stream-Table Join(流加维度表)
CREATE STREAM enriched_clicks AS
SELECT
    c.user_id,
    c.page,
    u.country
FROM clicks c
LEFT JOIN users u
  ON c.user_id = u.user_id
EMIT CHANGES;

-- 4) 30 分钟 Session 窗口(用户活跃 session)
SELECT
    user_id,
    COUNT(*) AS events,
    WINDOWSTART AS session_start,
    WINDOWEND   AS session_end
FROM clicks
WINDOW SESSION (30 MINUTES)
GROUP BY user_id
EMIT CHANGES;

EMIT CHANGES 表示「持续输出变更」,这是 ksqlDB 流处理的标志。EMIT FINAL 则只输出窗口关闭时的最终值。

10.4 ksqlDB 的优势

  • 数据分析师 / SQL 用户上手 5 分钟;
  • 拉起一个 docker 即可,无需写 Java;
  • 默认 connector 集成(CREATE SOURCE CONNECTOR ... WITH (...));
  • pull query(SELECT ... WHERE user_id='1')能像 KV 数据库一样查 KTable。

10.5 ksqlDB 的局限

  • 不支持复杂的事件时间逻辑 / CEP;
  • 没法写 UDAF(用户自定义聚合函数)以外的复杂状态机;
  • 真正的「重业务流处理」还是要 Streams DSL 或 Flink。

维度Kafka StreamsApache Flink
部署形态普通 Java app独立集群
调度JobManager / TaskManager
Source仅 Kafka30+ Connector(Kafka / MySQL / 文件 / Pulsar / Kinesis…)
状态后端RocksDB + Kafka TopicRocksDB + DFS Snapshot
Exactly OnceEOS v2(仅 Kafka)Two-Phase Commit(多 Sink)
CEPFlink CEP
窗口Tumbling/Hopping/Sliding/Session+ Custom WindowAssigner
SQLksqlDB(独立)Flink SQL(集成)
学习曲线平缓
运维成本低(无需独立集群)
适合场景Kafka-only 微服务化流处理多源、复杂、企业级

实战建议

  • 90% 的「实时统计 / 实时大屏 / 实时风控」用 Kafka Streams 或 ksqlDB 就够;
  • 上 Flink 的理由通常是「需要接非 Kafka 数据源」、「跑 Flink CEP / SQL 已经成基础设施」、「跨多个团队共享同一个 Flink 集群」。

12. 生产踩坑

坑 1:忘了配 application.id 唯一

application.id 同时是 consumer group + 内部 changelog topic 前缀。两个不同的应用用同一个 application.id → state 互相覆盖,灾难现场。

坑 2:num.stream.threads 设过大

线程数 > 总 Task 数,多余线程空转浪费 RAM。配置时让 总实例数 × num.stream.threads ≈ Topic 总分区数 即可。

坑 3:State Store 撑爆磁盘

KTable / 窗口聚合的状态默认存 RocksDB(本地磁盘)。窗口太长 / key 基数太大 → 磁盘爆。 解决:

  • 缩短 window;
  • 调 RocksDB block cache(rocksdb.config.setter);
  • KStream + 外部 KV(Redis)替代。

坑 4:repartition Topic 副本数 1

auto.create.topics.enable 为 true 时,Streams 内部创建的 repartition / changelog Topic 副本数 = 1,broker 重启可能丢数据。 解决:手动 streams.cleanUp() 后用 --replication-factor 3 创建,或配 replication.factor=3

坑 5:EOS 开了但 commit.interval.ms 太大

EOS commit interval 决定可见性延迟(消息要等 commit 后下游才看得见)。默认 100ms,一些读者把它调成 30s 后抱怨「下游延迟太大」。

坑 6:Stream-Stream Join 没考虑 grace period

JoinWindows 默认 grace 24h(旧版)或 0(新版),新版默认 0 会丢迟到消息;旧版默认 24h 会让 State Store 累积大量历史。显式指定 grace 是好习惯。

坑 7:不同实例的状态规模不均

key 倾斜(如热门用户 ID)→ 一个 Task 状态特别大、其他空闲。需要在业务层散列 key(加随机后缀)或换聚合维度。

坑 8:直接重启实例改了 application.id

application.id 一变 = 全新应用 = 重新建 Changelog + 重读 source topic。线上千万别手抖改。


13. 小结

  • Kafka Streams 是一个 Java 库(非集群),把流处理嵌进微服务,水平扩展靠加副本。
  • 流-表二元性:KStream 是 binlog,KTable 是当前快照,二者可无损转换。
  • 三种抽象:KStream(事件)/ KTable(状态)/ GlobalKTable(全量小维表)。
  • 算子分两类:无状态(mapfilterflatMapbranchpeek)和有状态(groupByKeyaggregatereducecountjoinwindow)。
  • State Store = RocksDB + Changelog Topic(compacted),状态本质上是一个 Kafka Topic,因此「状态可恢复」是天生的。
  • 三种 Join:KStream-KStream(必须 window)、KStream-KTable(流来一条查一次)、KTable-KTable(双表共维护)。
  • 四种窗口:Tumbling / Hopping / Sliding / Session,配合时间语义(Event/Processing/Ingestion)+ grace period 处理乱序。
  • EOS v2:一行 processing.guarantee=exactly_once_v2,把消费 + 写 + State 打包事务化。
  • ksqlDB = Streams 的 SQL 外壳,用 SQL 描述流处理。CREATE STREAM/TABLE + WINDOW TUMBLING 即可。
  • 与 Flink 选型:Kafka-only + 不想多维护集群 → Streams;多源 + 复杂 + 已有 Flink 基础设施 → Flink。

🎯 面试高频题

Q1:什么是流-表二元性?KStream / KTable / GlobalKTable 有什么区别?

考察点:核心概念、模型理解。

答案

  1. 流-表二元性:流(KStream)是事件序列(每条都是「发生了一件事」),表(KTable)是按 key 的当前状态(最新值覆盖旧值)。流可以聚合成表(重放变更),表可以发出变更流(Changelog)。它们是同一信息的两种表达:流是电影,表是截图。
  2. KStream:每条消息独立,重复 key 不会覆盖;按 key 分区;适合事件、点击流、订单流。
  3. KTable:按 key 取最新值;按 key 分区;底层是 RocksDB + Changelog Compacted Topic;适合「当前状态」、累加结果、维度表。
  4. GlobalKTable:和 KTable 一样按 key 取最新值,但 每个实例都全量复制。适合小维度表(万级行),可以无视分区与流 join,是「广播 join」。
  5. 应用差别
    • 状态聚合 / 当前余额 → KTable;
    • 事件 / 点击 → KStream;
    • 国家代码、商品类目这种小维表 → GlobalKTable;
  6. 加分项:提到 KTable 与流 Join 必须共分区,GlobalKTable 不需要;提到 GlobalKTable 不参与 Stream Task 并行(每实例 1 份),所以不能太大。

Q2:State Store 是怎么做到「状态可恢复」的?

考察点:底层存储、容错。

答案

  1. 本地存储:默认 RocksDB(key-value DB),存在每个 Streams 实例的本地磁盘;读写都是毫秒级。
  2. Changelog Topic:每次 State Store 更新(put / delete),Streams 同步写一条到 KAFKA-STREAMS-{app-id}-{store}-changelog Topic(compaction 模式),保证「磁盘 + Kafka」两份。
  3. 故障恢复:实例宕机 / 迁移 → 新实例从 Changelog 重放,把 RocksDB 重建。状态本质上就是一个 compacted topic
  4. standby replicasnum.standby.replicas=1 让另一个实例热备一份 State Store,主挂秒级接管。
  5. Compacted Topic 的角色:保证「同 key 永远只保留最新值」,Changelog 体积可控;触发 Compaction 后旧值被回收。
  6. 对比 Flink:Flink 也用 RocksDB + checkpoint 到 DFS(HDFS / S3);Streams 没有 DFS 依赖,状态完全靠 Kafka,运维更轻。
  7. 加分项:提到 changelog 的 cleanup.policy=compact、副本数应 ≥ 2;提到 standby replicas 还能服务 Interactive Queries(IQ)。

Q3:四种窗口(Tumbling / Hopping / Sliding / Session)分别是什么?给典型场景。

考察点:窗口类型、应用场景。

答案

  1. Tumbling(滚动):固定大小、不重叠、对齐时间。例:「每 1 分钟的 GMV」「每天的活跃用户数」。
  2. Hopping(跳跃 / 滑动):固定大小 + step;可重叠;对齐。例:「每 10 秒输出一次最近 1 分钟的实时 QPS 曲线」(窗口 1min,step 10s,相邻窗口重叠 50s)。
  3. Sliding(事件级滑动):每来一个事件就计算一次最近 N 时间内的状态;不对齐。例:「反欺诈 - 最近 30 秒同一用户交易笔数 > 5 报警」。
  4. Session:动态长度,由「无活动 gap」决定;相邻事件间隔超过 gap 就开新会话。例:「用户活跃 session(30 分钟无点击切断)」「设备上下线段统计」。
  5. 类比:考勤表里的不同打卡周期 —— Tumbling 像「每天 9:00-18:00」,Hopping 像「每小时统计最近 4 小时」,Session 像「员工进出门」,Sliding 像「每次打卡看最近 30 分钟」。
  6. 代码区别
    java
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)));
    .windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)).advanceBy(Duration.ofSeconds(10)));
    .windowedBy(SlidingWindows.ofTimeDifferenceWithNoGrace(Duration.ofSeconds(30)));
    .windowedBy(SessionWindows.ofInactivityGapWithNoGrace(Duration.ofMinutes(30)));
  7. 加分项:提到 grace period 处理迟到消息;提到 EMIT CHANGES(每次更新)vs EMIT FINAL(窗口关闭再发)。

Q4:KStream-KStream Join、KStream-KTable Join、KTable-KTable Join 区别?

考察点:Join 模型、共分区。

答案

  1. KStream-KStream:双流匹配,必须有 window(否则 State 无限增长)。例:订单流 ⨝ 支付流(10 分钟内匹配)。
  2. KStream-KTable:每来一条流消息查一次 KTable 当前值,不需要 window。例:订单流加用户档案。流来才触发 join,KTable 更新不触发。
  3. KTable-KTable:双表 join,不需要 window;任一边变化都触发 join。例:用户表 ⨝ 等级表。
  4. 共分区要求:以上三种都要求双方按相同 key 分区(不一致会自动 repartition,性能损耗)。唯一例外:与 GlobalKTable join 不要求共分区(小维表广播)。
  5. GlobalKTable 用途:当 join 的右侧是「小且全」(如国家代码、汇率表)时,用 GlobalKTable 避免 repartition + 全量在每实例本地查。
  6. 代码示例
    java
    // S-S window join
    ordersStream.join(paymentsStream, (o,p) -> ..., JoinWindows.of(Duration.ofMinutes(10)));
    // S-T join
    ordersStream.join(usersTable, (o,u) -> ...);
    // T-T join
    usersTable.join(levelsTable, (u,l) -> ...);
  7. 加分项:提到 inner / left / outer join;提到 KStream-KTable 是 Streams 里最常用、性能最好的 join;提到 KTable-KTable foreign key join(2.4+)。

考察点:EOS 协议、流处理语义。

答案

  1. 开关processing.guarantee=exactly_once_v2(Streams 2.6+)。
  2. 实现机制
    • 消费 source topic + 写 sink topic + 更新 State Store + 写 changelog 打包成一个事务 Producer 事务
    • sendOffsetsToTransaction 把 source consumer offset 也写到事务里;
    • commit 时一次性原子完成;
    • 中途崩溃 → 整个事务 abort,下次从未提交的 offset 重做。
  3. 下游配合:消费者必须 isolation.level=read_committed,否则会读到中间未提交消息。
  4. v1 vs v2:v1 是「每个 Task 一个事务 Producer」,资源开销大、吞吐损耗 ~30%;v2 改成「每个实例一个事务 Producer」(Producer-per-instance),损耗降到 ~10%。
  5. 与 Flink EOS 区别
    • Streams:仅在 Kafka 内部生效(source、sink、state 都在 Kafka 上);
    • Flink:基于 Two-Phase Commit Sink,可对接非 Kafka 系统(外部数据库);
    • Streams 实现简单(依赖 Kafka 事务),Flink 实现通用(依赖 sink 的 2PC 协议)。
  6. 加分项:提到 commit.interval.ms 越小延迟越低但吞吐越低(默认 100ms);提到 EOS 不能解决「调用外部 API 副作用」的幂等(只在 Kafka topic + State 之间生效)。

Q6:Kafka Streams 是怎么水平扩展的?Task 和 Partition 是什么关系?

考察点:调度模型、Rebalance。

答案

  1. 关系
    • 一个 Topology 被切成多个 Sub-Topology(按 repartition 划分);
    • 每个 Sub-Topology × 每个分区 = 1 个 Stream Task
    • Task 是最小调度单元,一个 Task 由一个线程独占
  2. 水平扩展:加 Streams 实例 → 通过 Consumer Group 协议 Rebalance → Task 自动迁移到新实例(含 State Store 重建 / 从 standby 接管)。
  3. 线程数num.stream.threads 决定每实例多少线程;总线程 = 实例数 × 每实例线程数;线程数 ≤ Task 数才不浪费。
  4. 状态迁移:Task 迁到新实例时,新实例从 changelog 重放重建 State Store;如果有 standby replica,秒级接管。
  5. Rebalance 协议:Cooperative Sticky(默认 2.4+)—— 增量再平衡,避免 Stop-the-World。
  6. 类比 Flink:Flink 的 Slot ≈ Streams 的 Task;Flink 的 TaskManager ≈ Streams 实例;只是 Streams 没有 JobManager / 调度器,Rebalance 由 Kafka 协议自己搞定。
  7. 加分项:提到 Rebalance 听 KafkaStreams.setUncaughtExceptionHandler / StateRestoreListener 监听状态恢复进度;提到分区数决定并行上限,Topic 分区数选少了 = 上限低。

本章配套:

  • 17_streams/demo.html:流-表二元性 + 4 种窗口 + Stream-Table Join + WordCount Topology 可视化。
  • 17_streams/code/wordcount_streams.java:经典 WordCount Java 版。
  • 17_streams/code/wordcount_faust.py:Python Faust 等价实现。
  • 17_streams/code/realtime_topn.py:实时 Top-N 排行榜(Faust)。
  • 17_streams/code/windowed_aggregate.py:1 分钟滚动窗口聚合(Faust)。
  • 17_streams/code/ksql_examples.sql:ksqlDB 常见语句集合。

🔗 延伸阅读

  • 第 11 章 消费者组与 Rebalance —— Streams 实例之间的 Cooperative Rebalance 是同一个协议。
  • 第 13 章 幂等与事务 —— EOS v2 复用的事务 Producer + read_committed 协议。
  • 第 14 章 Log Compaction —— Changelog Topic 用的就是 Compaction。
  • 第 16 章 Kafka Connect —— Streams 处理后的结果可以用 Sink Connector 落地。
  • 第 18 章 Schema Registry —— Streams 用 Avro Serdes 时依赖。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

sql
-- ====================================================================
-- 第 17 章 - Kafka Streams 与 ksqlDB
-- ksql_examples.sql - ksqlDB 常见语句集合
-- --------------------------------------------------------------------
-- 适用版本:ksqlDB 0.29.0+
-- 启动 ksqlDB CLI:
--   docker exec -it ksqldb-cli ksql http://ksqldb-server:8088
-- ====================================================================

-- ===== 0. 集群信息 =====
SHOW PROPERTIES;
SHOW STREAMS;
SHOW TABLES;
SHOW QUERIES;
SHOW TOPICS;
SHOW CONNECTORS;

-- ====================================================================
-- 1. 创建 STREAM(KSTREAM)= 把 Kafka Topic 当作事件流暴露
-- --------------------------------------------------------------------
-- 关键点:
--   * KEY 字段会作为消息 key
--   * TIMESTAMP 指定 event-time 来源(也可省略用 ROWTIME)
--   * VALUE_FORMAT:JSON / AVRO / PROTOBUF / DELIMITED
-- ====================================================================
CREATE STREAM clicks (
    user_id  VARCHAR KEY,
    page     VARCHAR,
    referer  VARCHAR,
    ts       BIGINT
) WITH (
    KAFKA_TOPIC   = 'learn.clicks',
    VALUE_FORMAT  = 'JSON',
    TIMESTAMP     = 'ts',
    PARTITIONS    = 6,
    REPLICAS      = 3
);


-- ====================================================================
-- 2. 创建 TABLE(KTABLE)= 按 key 取最新值
-- --------------------------------------------------------------------
-- 表的 source topic 必须是 compacted(或者你接受历史值被淘汰)
-- ====================================================================
CREATE TABLE users (
    user_id  VARCHAR PRIMARY KEY,
    name     VARCHAR,
    country  VARCHAR,
    vip      BOOLEAN
) WITH (
    KAFKA_TOPIC  = 'learn.users',
    VALUE_FORMAT = 'JSON'
);


-- ====================================================================
-- 3. 持续查询(Push Query):流处理的灵魂
-- --------------------------------------------------------------------
-- EMIT CHANGES 表示「持续输出变更」,与传统 SQL 一次性返回不同
-- ====================================================================

-- 3.1 简单过滤 + 投影
CREATE STREAM tech_clicks AS
    SELECT user_id, page
    FROM clicks
    WHERE page LIKE '/tech/%'
    EMIT CHANGES;

-- 3.2 改 key(会触发内部 repartition)
CREATE STREAM clicks_by_page AS
    SELECT page AS PAGE_KEY, user_id, ts
    FROM clicks
    PARTITION BY page
    EMIT CHANGES;


-- ====================================================================
-- 4. 窗口聚合:4 种窗口
-- ====================================================================

-- 4.1 Tumbling 1 分钟:每分钟 PV
CREATE TABLE pv_per_minute AS
    SELECT
        page,
        COUNT(*)                  AS pv,
        COUNT_DISTINCT(user_id)   AS uv
    FROM clicks
    WINDOW TUMBLING (SIZE 1 MINUTE)
    GROUP BY page
    EMIT CHANGES;

-- 4.2 Hopping:每 10s 输出一次最近 1min 的 QPS
CREATE TABLE pv_hopping AS
    SELECT page, COUNT(*) AS pv_1m
    FROM clicks
    WINDOW HOPPING (SIZE 1 MINUTE, ADVANCE BY 10 SECONDS)
    GROUP BY page
    EMIT CHANGES;

-- 4.3 Session 窗口:用户活跃 session(30 分钟无操作切断)
CREATE TABLE user_sessions AS
    SELECT
        user_id,
        COUNT(*)                                  AS events,
        TIMESTAMPTOSTRING(WINDOWSTART, 'HH:mm')   AS session_start,
        TIMESTAMPTOSTRING(WINDOWEND,   'HH:mm')   AS session_end
    FROM clicks
    WINDOW SESSION (30 MINUTES)
    GROUP BY user_id
    EMIT CHANGES;


-- ====================================================================
-- 5. JOIN
-- ====================================================================

-- 5.1 Stream-Table Join:流 ⨝ 维度表(最常用、性能最好)
CREATE STREAM enriched_clicks AS
    SELECT
        c.user_id,
        c.page,
        u.name      AS user_name,
        u.country   AS country,
        u.vip       AS is_vip
    FROM clicks c
    LEFT JOIN users u
      ON c.user_id = u.user_id
    EMIT CHANGES;

-- 5.2 Stream-Stream Join:双流匹配(必须 WITHIN)
-- 例:订单流 ⨝ 支付流(10 分钟内匹配)
CREATE STREAM orders_with_payment AS
    SELECT
        o.order_id,
        o.amount,
        p.tx_id     AS payment_tx,
        p.paid_at
    FROM orders o
    INNER JOIN payments p
      WITHIN 10 MINUTES
      ON o.order_id = p.order_id
    EMIT CHANGES;

-- 5.3 Table-Table Join
CREATE TABLE user_with_level AS
    SELECT
        u.user_id,
        u.name,
        l.level
    FROM users u
    JOIN levels l
      ON u.user_id = l.user_id
    EMIT CHANGES;


-- ====================================================================
-- 6. Pull Query:像 KV 数据库一样查 KTable
-- --------------------------------------------------------------------
-- ⚠ 仅 TABLE 支持 pull query,且查询字段必须是 key
-- ====================================================================
SELECT * FROM users WHERE user_id = 'u-1';

-- 窗口聚合表:必须指定 WINDOW 范围
SELECT * FROM pv_per_minute
WHERE page = '/tech/article-1'
  AND WINDOWSTART > 1700000000000;


-- ====================================================================
-- 7. CREATE STREAM AS SELECT 与持久查询管理
-- ====================================================================
-- 上面所有 CREATE STREAM/TABLE AS SELECT 都会生成一个「持久 Push Query」
SHOW QUERIES;             -- 列出
EXPLAIN <query_id>;       -- 看拓扑
TERMINATE <query_id>;     -- 停止
DROP STREAM tech_clicks;  -- 删除(必须先 TERMINATE 依赖)


-- ====================================================================
-- 8. UDF / UDAF(用户自定义函数)
-- --------------------------------------------------------------------
-- ksqlDB 内置丰富 UDF:UCASE / LCASE / SPLIT / REGEXP_EXTRACT / CASE WHEN
-- 时间:TIMESTAMPTOSTRING / STRINGTOTIMESTAMP / DATEADD / EXTRACT
-- JSON:EXTRACTJSONFIELD / JSON_OBJECT / TO_JSON_STRING
-- 集合:COLLECT_LIST / COLLECT_SET / EARLIEST_BY_OFFSET / LATEST_BY_OFFSET
-- --------------------------------------------------------------------
SELECT
    user_id,
    LATEST_BY_OFFSET(page)            AS last_page,
    COLLECT_LIST(page)                AS visited_pages
FROM clicks
WINDOW TUMBLING (SIZE 5 MINUTES)
GROUP BY user_id
EMIT CHANGES;


-- ====================================================================
-- 9. Connector 集成(直接在 ksqlDB 里管理 Kafka Connect)
-- ====================================================================
CREATE SOURCE CONNECTOR debezium_orders WITH (
    'connector.class'      = 'io.debezium.connector.mysql.MySqlConnector',
    'database.hostname'    = 'mysql',
    'database.port'        = '3306',
    'database.user'        = 'debezium',
    'database.password'    = 'debezium-secret',
    'database.server.id'   = '1234',
    'topic.prefix'         = 'dbz.shop',
    'table.include.list'   = 'shop.orders'
);

SHOW CONNECTORS;
DESCRIBE CONNECTOR debezium_orders;
DROP CONNECTOR debezium_orders;


-- ====================================================================
-- 10. INSERT INTO(手动塞数据,调试方便)
-- ====================================================================
INSERT INTO clicks (user_id, page, ts)
VALUES ('u-1', '/tech/kafka', UNIX_TIMESTAMP());


-- ====================================================================
-- 11. SET 配置(仅本会话生效)
-- ====================================================================
-- 从最早开始消费
SET 'auto.offset.reset' = 'earliest';
-- 显式指定 EOS
SET 'processing.guarantee' = 'exactly_once_v2';
-- 增大并行度
SET 'num.stream.threads' = '4';


-- ====================================================================
-- 调试技巧:
--   * EXPLAIN <query_id>;     看 Streams Topology
--   * PRINT 'topic_name' FROM BEGINNING LIMIT 10;  打印 topic 内容
--   * DESCRIBE EXTENDED users;  看 schema、统计、所有依赖
-- ====================================================================
python
"""
第 17 章 - Kafka Streams 与 ksqlDB
realtime_topn.py - 用 Faust 实现「实时 Top-N 商品销量排行榜」

业务场景:
    电商大屏需要「最近 1 小时下单数 Top-10 商品」。每个订单是一个事件,
    发到 `learn.orders` Topic;Streams 应用维护一个滑动窗口聚合 + Top-N
    Heap,定时把 Top-10 推到 `learn.topn` Topic 给前端 BFF 拉取。

关键技巧:
    1) 1 小时滚动窗口聚合每商品销量(KTable<product_id, count>)
    2) 维护一个全局 Top-N Heap(Table,单 key='topn')
    3) 每秒输出一次最新 Top-10 到下游 Topic / Web 接口

依赖:
    pip install faust-streaming python-rocksdb

准备:
    kafka-topics.sh ... --create --topic learn.orders --partitions 6 --replication-factor 1
    kafka-topics.sh ... --create --topic learn.topn   --partitions 1 --replication-factor 1

运行:
    faust -A realtime_topn worker -l info --web-port 6068
    curl http://localhost:6068/topn/now
"""
from __future__ import annotations

import faust
import heapq
from datetime import timedelta
from typing import List, Tuple

WINDOW_SIZE = timedelta(hours=1)
EMIT_EVERY  = timedelta(seconds=5)
TOP_N       = 10


class Order(faust.Record, serializer="json"):
    order_id:   str
    product_id: str
    user_id:    str
    amount:     float
    ts:         float    # event time, seconds


class TopNEntry(faust.Record, serializer="json"):
    product_id: str
    count:      int


app = faust.App(
    "realtime-topn",
    broker="kafka://localhost:9092",
    store="rocksdb://",
    table_standby_replicas=1,
)

orders_topic = app.topic("learn.orders", value_type=Order)
topn_topic   = app.topic("learn.topn",   value_type=List[TopNEntry])

# 1 小时滚动窗口聚合:每个 product_id 的下单笔数
hourly_count = (
    app.Table("hourly-count", default=int, partitions=6)
       .tumbling(WINDOW_SIZE, expires=timedelta(hours=2))
       .relative_to_field(Order.ts)
)

# 用一个 single-key Table 存当前 Top-N 快照(方便 IQ 查询)
topn_snapshot = app.Table("topn-snapshot", default=list, partitions=1)


# -------------------------------------------------------------------
# Agent 1:消费订单流,更新窗口聚合
# -------------------------------------------------------------------
@app.agent(orders_topic)
async def aggregate(stream):
    """每来一个订单:对应 product_id 的窗口计数 +1。"""
    async for order in stream:
        # Faust 的 windowed Table:直接用 [] 操作即可
        hourly_count[order.product_id] += 1


# -------------------------------------------------------------------
# Timer:定期扫描全表,算出 Top-N,写到 snapshot + Kafka
# -------------------------------------------------------------------
@app.timer(interval=EMIT_EVERY.total_seconds())
async def emit_topn():
    """周期性把当前窗口的 Top-N 推送到下游 + 内存快照。"""
    # 当前窗口下的所有 (product, count)
    # 注意:Faust 的 windowed Table 遍历需要 .items() + .current()
    pairs: List[Tuple[str, int]] = []
    for k, w in hourly_count.items():
        try:
            cnt = w.current()
            if cnt > 0:
                pairs.append((k, cnt))
        except Exception:
            continue

    # heapq.nlargest = Top-N
    top = heapq.nlargest(TOP_N, pairs, key=lambda x: x[1])
    payload = [TopNEntry(product_id=p, count=c) for p, c in top]

    # 写到 snapshot Table(IQ 查询)
    topn_snapshot["current"] = [(e.product_id, e.count) for e in payload]
    # 写到下游 Topic
    await topn_topic.send(value=payload)

    if payload:
        print(f"[Top-{TOP_N}] " +
              " | ".join(f"{e.product_id}={e.count}" for e in payload))


# -------------------------------------------------------------------
# Web 接口:实时返回 Top-N(Interactive Query)
# -------------------------------------------------------------------
@app.page("/topn/now")
async def http_topn(web, request):
    items = topn_snapshot.get("current", [])
    return web.json([{"product_id": p, "count": c} for p, c in items])


if __name__ == "__main__":
    app.main()


# ============================================================================
# 性能与扩展
# ----------------------------------------------------------------------------
# - 商品总数 1M+ 时,timer 里全表扫描会变慢,建议改成「维护 Top-N 增量更新」
#   (比如只在 cnt 变化时检查是否能替换 heap 末尾)。
# - 窗口越长,State Store 越大;可以加 `compacting` 后端 + 缩短 expires。
# - 真正高吞吐(百万 QPS)建议用 Java Streams 或 Flink;Faust 里 Python
#   单进程吞吐有 GIL 上限。
# ============================================================================
python
"""
第 17 章 - Kafka Streams 与 ksqlDB
windowed_aggregate.py - 用 Faust 演示「1 分钟滚动窗口聚合 GMV」

业务场景:
    实时大屏需要「按分钟聚合的订单总额(GMV)」,按支付状态分组,
    输出到下游 `learn.gmv-per-minute` Topic 给前端拉取。

核心知识点:
    1) Tumbling Window(1 分钟,不重叠)
    2) Aggregate(不仅 count,还累加 amount)
    3) Event Time(用消息里的 ts 而不是处理时间)
    4) Grace Period 处理乱序消息

依赖:
    pip install faust-streaming python-rocksdb

准备:
    kafka-topics.sh ... --create --topic learn.orders          --partitions 3
    kafka-topics.sh ... --create --topic learn.gmv-per-minute  --partitions 1

运行:
    faust -A windowed_aggregate worker -l info --web-port 6069
"""
from __future__ import annotations

import faust
from datetime import timedelta
from typing import NamedTuple

WINDOW_SIZE  = timedelta(minutes=1)    # 1 分钟滚动
GRACE_PERIOD = timedelta(seconds=30)   # 允许 30s 迟到


class Order(faust.Record, serializer="json"):
    order_id: str
    user_id:  str
    amount:   float
    status:   str           # PENDING / PAID / CANCELLED / REFUNDED
    ts:       float         # event time (epoch seconds)


class MinuteStats(faust.Record, serializer="json"):
    bucket_minute: int      # epoch minute
    status:        str
    order_cnt:     int
    gmv:           float
    avg_amount:    float


app = faust.App(
    "windowed-gmv",
    broker="kafka://localhost:9092",
    store="rocksdb://",
)

orders_topic = app.topic("learn.orders", value_type=Order)
gmv_topic    = app.topic("learn.gmv-per-minute", value_type=MinuteStats)


# -------------------------------------------------------------------
# 一个 windowed Table,key = status,value = (count, sum)
# -------------------------------------------------------------------
class Acc(NamedTuple):
    cnt: int
    sum: float


minute_acc = (
    app.Table("minute-acc",
              default=lambda: Acc(0, 0.0),
              partitions=3)
       .tumbling(WINDOW_SIZE, expires=timedelta(hours=1))
       .relative_to_field(Order.ts)   # 按 event time 做窗口
)


# -------------------------------------------------------------------
# Agent:每来一条订单,按 status 累加 cnt + sum
# -------------------------------------------------------------------
@app.agent(orders_topic)
async def consume_orders(stream):
    async for order in stream:
        cur = minute_acc[order.status].value()
        minute_acc[order.status] = Acc(cur.cnt + 1, cur.sum + order.amount)


# -------------------------------------------------------------------
# 周期性扫描当前窗口,输出聚合
# -------------------------------------------------------------------
@app.timer(interval=10.0)
async def flush_to_topic():
    """每 10 秒把当前窗口的快照推到下游 Topic。"""
    for status, w in minute_acc.items():
        try:
            acc: Acc = w.current()
        except Exception:
            continue
        if acc.cnt == 0:
            continue
        # current() 返回当前 wall-clock 落入的窗口聚合
        # 拿当前的 minute bucket
        import time
        bucket = int(time.time() // 60)
        stats = MinuteStats(
            bucket_minute=bucket,
            status=status,
            order_cnt=acc.cnt,
            gmv=round(acc.sum, 2),
            avg_amount=round(acc.sum / acc.cnt, 2),
        )
        await gmv_topic.send(key=f"{bucket}:{status}", value=stats)
        print(f"[GMV] minute={bucket} status={status} "
              f"cnt={stats.order_cnt} gmv={stats.gmv} avg={stats.avg_amount}")


# -------------------------------------------------------------------
# 实时查询接口:当前窗口的实时数字
# -------------------------------------------------------------------
@app.page("/gmv/current")
async def http_current(web, request):
    out = []
    for status, w in minute_acc.items():
        try:
            acc = w.current()
            if acc.cnt > 0:
                out.append({"status": status, "cnt": acc.cnt,
                            "gmv": round(acc.sum, 2)})
        except Exception:
            pass
    return web.json(out)


if __name__ == "__main__":
    app.main()


# ============================================================================
# 等价的 Java Streams DSL 写法(仅参考,不在本文件中运行)
# ----------------------------------------------------------------------------
# StreamsBuilder builder = new StreamsBuilder();
# KStream<String, Order> orders = builder.stream("learn.orders",
#     Consumed.with(Serdes.String(), orderSerde)
#             .withTimestampExtractor(new OrderTsExtractor()));
#
# KTable<Windowed<String>, Acc> agg = orders
#     .groupBy((k, v) -> v.status)
#     .windowedBy(TimeWindows.ofSizeAndGrace(
#         Duration.ofMinutes(1),
#         Duration.ofSeconds(30)))
#     .aggregate(
#         () -> new Acc(0, 0.0),
#         (k, order, acc) -> new Acc(acc.cnt + 1, acc.sum + order.amount),
#         Materialized.<String, Acc, WindowStore<Bytes,byte[]>>as("minute-agg")
#                     .withValueSerde(accSerde));
#
# agg.toStream()
#    .map((wk, acc) -> KeyValue.pair(
#        wk.window().start() / 60_000 + ":" + wk.key(),
#        new MinuteStats(...)))
#    .to("learn.gmv-per-minute", Produced.with(Serdes.String(), statsSerde));
# ============================================================================
python
"""
第 17 章 - Kafka Streams 与 ksqlDB
wordcount_faust.py - 用 Python `faust-streaming` 实现的 WordCount

⚠ 注意:
    Kafka Streams 是 JVM 库,Python 没有官方等价物。`faust-streaming` 是社区在
    Robinhood `faust` 基础上 fork 出来的活跃版本,提供 Streams 风格的 API(agent +
    Table),底层用 RocksDB + Changelog Topic。

    Faust 的 Changelog 命名是 `<app-id>-<table>-changelog`,但格式与 Java Streams
    并不互通;不能在生产中混用。

    Faust 适合:Python 团队、轻量场景、原型快速验证;不推荐替代 JVM Streams
    跑高吞吐核心链路。

依赖:
    pip install faust-streaming python-rocksdb

准备 Kafka Topic:
    kafka-topics.sh --bootstrap-server localhost:9092 \\
      --create --topic learn.text --partitions 3 --replication-factor 1

运行:
    # Worker 进程(可启多个,自动按分区 Rebalance)
    faust -A wordcount_faust worker -l info --web-port 6066

    # 喂数据
    kafka-console-producer.sh --bootstrap-server localhost:9092 --topic learn.text
    > kafka streams kafka topic kafka

    # 查看 Table(实时排行)
    curl http://localhost:6066/wordcount/top
"""
from __future__ import annotations

import faust
import re
from collections import OrderedDict

# ---------- App / Topic / Table ----------
app = faust.App(
    "wordcount-py",                              # 等价 application.id
    broker="kafka://localhost:9092",
    value_serializer="raw",                      # 不要尝试 json 解码原始字符串
    store="rocksdb://",                          # State Store 后端
    table_standby_replicas=1,                    # 等价 num.standby.replicas
    topic_partitions=3,
    consumer_auto_offset_reset="earliest",
)

text_topic = app.topic("learn.text", value_type=bytes)

# Table 等价 KTable,落地到 RocksDB + Changelog Topic
counts = app.Table("counts-store", default=int, partitions=3)

WORD_RE = re.compile(r"\W+")


# ---------- Agent 等价 Streams DSL ----------
@app.agent(text_topic)
async def process(stream):
    """每收到一行文本:拆词 → 累加到 counts Table。"""
    async for line in stream:
        text = line.decode(errors="ignore").lower()
        for word in WORD_RE.split(text):
            if not word:
                continue
            counts[word] += 1
            print(f"[count] {word} -> {counts[word]}")


# ---------- HTTP 接口:Interactive Queries 风格 ----------
@app.page("/wordcount/{word}/")
async def get_word(web, request, word: str):
    """直接查 RocksDB 拿某个词的当前计数。"""
    return web.json({"word": word, "count": counts.get(word, 0)})


@app.page("/wordcount/top")
async def top_n(web, request):
    """返回 Top-10 高频词。"""
    items = sorted(counts.items(), key=lambda x: -x[1])[:10]
    return web.json([{"word": w, "count": c} for w, c in items])


if __name__ == "__main__":
    app.main()


# ============================================================================
# 与 Java Streams 对照
# ----------------------------------------------------------------------------
# | Java Streams                          | Faust                            |
# |---------------------------------------|----------------------------------|
# | StreamsBuilder.stream("learn.text")   | app.topic("learn.text")          |
# | KTable + Materialized.as("xxx")       | app.Table("xxx", default=int)    |
# | flatMapValues + groupBy + count       | for word in line.split: cnt[w]++ |
# | num.stream.threads                    | 启动多个 worker 进程             |
# | num.standby.replicas                  | table_standby_replicas           |
# | EOS v2                                | Faust 不支持端到端 EOS(仅 ALO) |
# | Interactive Queries                   | @app.page web 接口               |
# ----------------------------------------------------------------------------
# 性能对照(参考):
#   - 单 worker 吞吐:Faust ≈ 5-20K msg/s(Python GIL 限制)
#   - JVM Streams 单线程:50-200K msg/s
#   - 高吞吐用 Streams / Flink;Python 团队的中小流量用 Faust 完全够。
# ============================================================================
java
/*
 * ====================================================================
 * 第 17 章 - Kafka Streams 与 ksqlDB
 * wordcount_streams.java - 经典 WordCount 的 Kafka Streams Java 实现
 * --------------------------------------------------------------------
 * Maven 依赖(pom.xml):
 *
 *   <dependency>
 *     <groupId>org.apache.kafka</groupId>
 *     <artifactId>kafka-streams</artifactId>
 *     <version>3.8.0</version>
 *   </dependency>
 *
 * 准备:
 *   kafka-topics.sh --bootstrap-server localhost:9092 \
 *     --create --topic learn.text       --partitions 3 --replication-factor 1
 *   kafka-topics.sh --bootstrap-server localhost:9092 \
 *     --create --topic learn.wc-output  --partitions 3 --replication-factor 1
 *
 * 运行:
 *   mvn package
 *   java -jar target/wordcount-1.0.jar
 *
 * 输入:
 *   kafka-console-producer.sh --bootstrap-server localhost:9092 --topic learn.text
 *   > kafka streams kafka topic kafka
 *
 * 观察输出:
 *   kafka-console-consumer.sh --bootstrap-server localhost:9092 \
 *     --topic learn.wc-output --from-beginning \
 *     --property print.key=true \
 *     --value-deserializer org.apache.kafka.common.serialization.LongDeserializer
 * ====================================================================
 */
package com.example.kafka.streams;

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Materialized;
import org.apache.kafka.streams.kstream.Produced;

import java.util.Arrays;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class WordCount {

    public static void main(String[] args) {
        // ===== 1. 配置 =====
        Properties props = new Properties();
        // application.id 同时是 consumer group + 内部 changelog/repartition 的前缀
        // 必须全局唯一!
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        // 每个实例的线程数;total threads = instances × this
        props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 2);

        // EOS v2(可选):把消费 + 写 + State + changelog 打包成一个事务
        // 代价是吞吐 ~10% 损耗、commit interval 默认 100ms 的延迟
        // props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);

        // standby replica:另一实例热备一份 State Store,主挂秒级接管
        props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);

        // ===== 2. 构建 Topology =====
        StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> textLines = builder.stream("learn.text");

        KTable<String, Long> wordCounts = textLines
            // 一句话 → 多个 word(无状态算子)
            .flatMapValues(line -> Arrays.asList(line.toLowerCase().split("\\W+")))
            // 过滤掉空 word
            .filter((key, word) -> word != null && !word.isEmpty())
            // 改 key 为 word 自己;这一步会触发后续 repartition
            .groupBy((key, word) -> word)
            // 聚合:count,输出是 KTable<word, count>
            // Materialized.as 给 State Store 显式命名,便于 IQ 查询
            .count(Materialized.as("counts-store"));

        // 把 KTable 转回 KStream 并写到 sink topic
        wordCounts.toStream()
                  .to("learn.wc-output", Produced.with(Serdes.String(), Serdes.Long()));

        // ===== 3. 启动 =====
        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        CountDownLatch latch = new CountDownLatch(1);

        // 优雅关闭
        Runtime.getRuntime().addShutdownHook(new Thread("streams-shutdown") {
            @Override
            public void run() {
                streams.close();
                latch.countDown();
            }
        });

        // 状态恢复进度监听(生产建议接入)
        streams.setStateListener((newState, oldState) ->
            System.out.printf("[State] %s -> %s%n", oldState, newState));

        // Uncaught exception 处理:默认 SHUTDOWN_CLIENT;可改 REPLACE_THREAD(线程级容错)
        streams.setUncaughtExceptionHandler(ex -> {
            System.err.println("Uncaught: " + ex);
            return org.apache.kafka.streams.errors.StreamsUncaughtExceptionHandler
                       .StreamThreadExceptionResponse.REPLACE_THREAD;
        });

        // 打印 Topology(调试好帮手)
        System.out.println(builder.build().describe());

        try {
            streams.start();
            latch.await();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.exit(1);
        }
    }
}

/*
 * ====================================================================
 * 内部产生的 Kafka Topic(应用启动后会自动创建):
 *   wordcount-app-counts-store-changelog        ← KTable 的 Changelog(compacted)
 *   wordcount-app-KSTREAM-AGGREGATE-STATE-...   ← repartition topic(按 word 重分区)
 *
 * 想干净启动(清掉所有内部 topic + 本地 RocksDB):
 *   kafka-streams-application-reset.sh \
 *     --application-id wordcount-app \
 *     --input-topics learn.text \
 *     --bootstrap-server localhost:9092
 *   rm -rf /tmp/kafka-streams/wordcount-app
 *
 * ====================================================================
 */

ksql_examples.sql ↗ · realtime_topn.py ↗ · windowed_aggregate.py ↗ · wordcount_faust.py ↗ · wordcount_streams.java ↗