🌊 第 17 章 · Kafka Streams 与 ksqlDB 可视化

九个模块:① 流-表二元性 ② 4 种窗口对比 ③ KStream-KTable Join ④ WordCount Topology ⑤ KStream/KTable/GlobalKTable 选型 ⑥ 三种 Join 对比 ⑦ 时间语义 ⑧ State Store 容灾 ⑨ EOS v2 ⑩ ksqlDB 速查 ⑪ 常见坑。

① 流-表二元性(Stream-Table Duality)

左:KStream(事件流,每条都是一个事实);右:KTable(按 key 取最新值的当前状态)。点击「+ 加事件」感受两者关系。

KStream(learn.events)

KTable(balances)= sum by key

keybalanceupdates

② 四种窗口类型对比动画

点击上方按钮选择窗口类型。

③ KStream-KTable Join(流加维度表)

KTable 维护用户档案(按 user_id 取最新);KStream 进来订单时用 user_id 查 KTable 加上用户名 / 国家。

users (KTable)

usernamecountry

orders (KStream)

enriched (输出 KStream)

④ WordCount Topology 可视化

输入一句话,观察消息在拓扑节点之间流动。

📨
Source
learn.text
🔪
flatMapValues
split words
🔑
groupBy
by word
🔢
count
State Store
📤
Sink
learn.wc-output

State Store: counts-store

wordcount

Changelog Topic(compacted)

(等待输入…)

⑤ KStream / KTable / GlobalKTable 选型

抽象语义分区典型用途反例
KStream 事件流,每条都是独立事实 按消息 key 点击日志、订单事件、IoT 上报 不要存最新状态(每条都保留)
KTable 按 key 取最新值的实时快照 按 key(与 KStream 同分区策略) 用户档案、商品价格、配置项 不要做事件计数(应改 KStream → count)
GlobalKTable 全集群每个 task 都拥有完整副本 不分区(整张表广播) 小型维度表、字典(≤ 100 MB) 不要给百万行大表用,每个实例都要装一份

💡 KStream-KTable Join 要求两边同 key 同分区;KStream-GlobalKTable Join 不需要 co-partitioning,可按任意外键查找。

⑥ 三种 Join 对比

Join 类型是否要 window触发典型场景
KStream-KStream ✅ 必须(如 5min) 两边都来才匹配 下单 + 支付,5 分钟内匹配;点击 + 转化
KStream-KTable ❌ 不需要 左边来一条就查右表 订单流 + 用户档案表 → 订单加用户名
KTable-KTable ❌ 不需要 任何一边变化都触发 账户余额 + 利率表 → 实时利息
KStream-GlobalKTable ❌ 不需要 左边来一条查全局表(可用外键) 事件流 + 全量小型字典(国家码、币种)
// KStream-KTable join 示例(不需要 window)
KStream<String, Order> orders = builder.stream("orders");
KTable<String, User>  users  = builder.table("users");
orders.join(users, (order, user) -> enrich(order, user))
      .to("orders.enriched");
  

⑦ 三种时间语义

类型来源优点缺点
Event Time 消息体里的业务时间(如 ts 字段) 正确反映业务顺序,乱序也可处理 需要自定义 TimestampExtractor,要处理 watermark / late event
Ingestion Time Broker 接收消息时的 server timestamp 无需用户配置;接近事件时间 客户端积压时会失真
Processing Time Streams 处理消息那一刻的本地系统时间 实现简单,无依赖 重放历史数据时窗口完全错位
// 自定义 Event Time 抽取
public class JsonTsExtractor implements TimestampExtractor {
    public long extract(ConsumerRecord<Object,Object> r, long partitionTime) {
        long ts = parseJson((String) r.value()).getLong("event_ts");
        return ts > 0 ? ts : partitionTime;   // 兜底:用上一次的时间
    }
}
props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, JsonTsExtractor.class);
  

💡 配置 grace.ms 接受迟到事件,但越长状态越大;金融场景常 grace=5min,普通分析 grace=10s

⑧ State Store + Changelog Topic(容灾)

RocksDB 在本地磁盘存状态;每次 put 同步写入 changelog topic(compacted)。Task 迁移到新 Worker 时,从 changelog 重放重建本地 store。

[Task-0] -- put(k=alice, v=10) --> RocksDB (本地)
                                |
                                +--> producer.send(changelog-topic, k=alice, v=10)

崩溃后新 Worker 接手 Task-0:
  1. 创建空 RocksDB
  2. 从 changelog topic 头到尾 replay → 状态完全恢复
  3. 切换为正常处理消息

⚠ 优化:
  - changelog 是 compacted topic(log.cleanup.policy=compact),同 key 只保留最新
  - num.standby.replicas=1:另一台 Worker 实时跟随 changelog → 故障接管秒级
  - 关闭:StoreBuilder.withLoggingDisabled() → 只本地、不容灾(仅适合纯衍生数据)
  

⑨ EOS v2 在 Streams 中的开启

props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
          StreamsConfig.EXACTLY_ONCE_V2);   // ← 一行打开
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3);   // 内部 topic ≥ 3
  
版本差别
EOS v1(已废弃)每个 input partition 单独事务,开销大
EOS v2(推荐)每个 thread 共享一个 producer + 事务,性能接近 ALOS(At Least Once Semantic)

⚠ 下游消费者必须 isolation.level=read_committed,否则会看到 abort 的中间值。

⑩ ksqlDB 常见 SQL 速查

-- 创建 STREAM(基于 Topic)
CREATE STREAM clicks (user_id VARCHAR, page VARCHAR, ts BIGINT)
  WITH (KAFKA_TOPIC='learn.clicks', VALUE_FORMAT='JSON', TIMESTAMP='ts');

-- 创建 TABLE(按 key 取最新)
CREATE TABLE users (user_id VARCHAR PRIMARY KEY, name VARCHAR, country VARCHAR)
  WITH (KAFKA_TOPIC='learn.users', VALUE_FORMAT='JSON');

-- 滚动窗口聚合:每分钟每用户 PV
CREATE TABLE pv_per_min AS
  SELECT user_id, COUNT(*) AS pv
  FROM clicks
  WINDOW TUMBLING (SIZE 1 MINUTE)
  GROUP BY user_id
  EMIT CHANGES;

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

-- Pull Query(点查 TABLE 当前值)
SELECT * FROM users WHERE user_id = 'u-1';

-- 持续查询管理
SHOW QUERIES;
TERMINATE CSAS_PV_PER_MIN_3;
DROP TABLE pv_per_min DELETE TOPIC;
  

⑪ 常见坑速查

表现修复
Co-partitioning 不一致 KStream-KTable join 出现 NPE / 漏匹配 两边 partition 数必须相同;上游不一致就先 repartition()
用 KStream.count() 直接当库存 每次重启数都翻倍 应该用 KTable + 业务事件 = 增量更新;或定期快照
State Store 巨大但没启 standby Worker 故障切换要恢复几十分钟 num.standby.replicas=1 + 监控 changelog 大小
Window 不设 grace 迟到 1 秒的事件全丢 .grace(Duration.ofSeconds(10))
Application.id 改了 所有 changelog/repartition topic 重建,状态全失 application.id = 业务身份证,永远不要随意改
EOS v2 没改下游 下游看到 abort 的脏数据 所有下游 Consumer / ksqlDB / Connect Sink 全部 read_committed

⑫ Streams / Flink / Spark Streaming 选型

维度Kafka StreamsFlinkSpark Structured Streaming
部署形态普通 Java 应用(Jar)独立集群(JM/TM)Spark 集群(Driver/Executor)
外部数据源仅 Kafka多种 source/sink多种
状态后端RocksDB + changelogRocksDB / 内存 + checkpointHDFS checkpoint
EOSv2,开箱即用checkpoint + 2PC sinkv3 后支持
窗口Tumbling / Hopping / Session / Sliding同 + 自定义触发器,更灵活Tumbling / Sliding / Session
SQLksqlDB(独立服务)Flink SQL / Table APISpark SQL
定位Kafka 内嵌轻量级流处理纯流计算引擎,金融实时首选批流统一,离线为主

💡 选型口诀:
· 处理 100% 在 Kafka 内部、想简单部署 → Kafka Streams / ksqlDB
· 多源 + 复杂窗口 + 大状态 + 金融级 EOS → Flink
· 已有 Spark 离线团队、需要批流统一 → Spark Structured Streaming

⑬ Streams 关键调优参数

参数默认建议原因
num.stream.threads1= max(input partition 数)每线程跑若干 task,增大并行度
num.standby.replicas01状态可秒级切换
cache.max.bytes.buffering10MB50~200MB聚合 dedupe,减少下游 RPS
commit.interval.ms30s(EOS 100ms)EOS 200~500ms降低 EOS 模式吞吐损耗
processing.guaranteeat_least_onceexactly_once_v2(金融)EOS v2 性能损耗 < 10%
buffered.records.per.partition100010000+抗瞬时峰值
state.dir/tmp/kafka-streams独立 SSD 目录RocksDB 强 IO
📚 第 17 章配套:17_streams.md · code/wordcount_streams.java · code/wordcount_faust.py
code/realtime_topn.py · code/windowed_aggregate.py · code/ksql_examples.sql