KafkaNotes

第 23 章:Kafka 生态:Connect、Streams 与 ksqlDB

zjc 于 2026-01-23 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 的强大一半来自生态。本章介绍三个官方系工具:Connect 解决“数据进出 Kafka”,Streams 解决“流处理”,ksqlDB 让你用 SQL 写流式应用。

23.1 Kafka Connect

是什么

Connect 是一个数据集成框架,以插件方式运行 connector,无需写代码即可在 Kafka 与外部系统之间搬运数据:

Source Connector: 外部系统 -> Kafka
  JDBC / Debezium(CDC) / File / S3 ...

Sink Connector: Kafka -> 外部系统
  JDBC / Elasticsearch / S3 / ClickHouse ...

核心概念

概念 说明
Connector 逻辑任务定义(连接哪个库、读哪些表)
Task 实际并行工作的执行单元
Worker 运行 connector/task 的进程
Converter 消息格式转换(JSON/Avro/Protobuf)
Offset source 端进度(如表的 binlog 位点),存于内部 topic
REST API 管理 connector 的 HTTP 接口

两种模式

standalone:单进程,适合开发测试
distributed:多 worker,任务自动均衡与容错,生产必选

启动(distributed)

bin/connect-distributed.sh config/connect-distributed.properties

关键配置:

bootstrap.servers=localhost:9092
group.id=connect-cluster
offset.storage.topic=connect-offsets
config.storage.topic=connect-configs
status.storage.topic=connect-status
offset.storage.replication.factor=3
config.storage.replication.factor=3
status.storage.replication.factor=3
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter

用 REST 创建 JDBC Source

curl -X POST http://localhost:8083/connectors \
  -H 'Content-Type: application/json' -d '{
  "name": "jdbc-orders-source",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "tasks.max": "4",
    "connection.url": "jdbc:mysql://db:3306/shop",
    "connection.user": "reader",
    "connection.password": "***",
    "table.whitelist": "orders",
    "mode": "incrementing",
    "incrementing.column.name": "id",
    "topic.prefix": "cdc-"
  }
}'

管理命令:

curl localhost:8083/connectors                 # 列出
curl localhost:8083/connectors/jdbc-orders-source/status
curl -X DELETE localhost:8083/connectors/jdbc-orders-source

选型建议

23.2 Kafka Streams

是什么

Streams 是一个 Java 库(不是独立服务),嵌入应用即可做有状态流计算:

StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> lines = builder.stream("input-topic");

KTable<String, Long> counts = lines
        .flatMapValues(v -> Arrays.asList(v.toLowerCase().split("\\W+")))
        .groupBy((k, word) -> word)
        .count(Materialized.as("word-counts"));

counts.toStream().to("output-topic", Produced.with(Serdes.String(), Serdes.Long()));

KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

配置:

props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

application.id 同时充当消费组 ID 与状态存储前缀,同一应用的所有实例必须一致。

KStream vs KTable

抽象 语义 类比
KStream 事件流,每条都重要 append 日志
KTable 按 key 聚合的最新状态 可更新的表
GlobalKTable 全实例广播的表 小维度表

经典组合:订单流(KStream)join 用户表(KTable),实时补全维度。

窗口

clicks.groupByKey()
      .windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(10)))
      .count();

状态与容错

状态存本地 RocksDB,并后台备份到 changelog topic

实例宕机 -> 新实例接管分区 -> 从 changelog 重建本地状态 -> 继续处理

因此 Streams 应用会自动创建 <app-id>-<store>-changelog 主题,不要手工删除。

适用场景

复杂拓扑、事件时间语义、大规模状态管理时,再评估 Flink(第 25 章)。

23.3 ksqlDB

ksqlDB 把流处理变成 SQL,适合快速验证与轻量场景。

-- 建流
CREATE STREAM orders (
  order_id VARCHAR KEY,
  user_id VARCHAR,
  amount DECIMAL(12, 2)
) WITH (
  KAFKA_TOPIC = 'order-events',
  VALUE_FORMAT = 'JSON'
);

-- 每 5 分钟统计各用户消费
CREATE TABLE user_spend AS
  SELECT user_id,
         COUNT(*) AS order_cnt,
         SUM(amount) AS total
  FROM orders
  WINDOW TUMBLING (SIZE 5 MINUTES)
  GROUP BY user_id
  EMIT CHANGES;

-- 持续查询(推模式)
SELECT * FROM user_spend
  WHERE total > 10000
  EMIT CHANGES;

优点:上手快、声明式、内置服务化。注意:复杂逻辑、精细状态控制与大型作业,仍建议 Streams/Flink。

23.4 选型对比

需求 推荐
数据库/日志/对象存储进出 Kafka Connect
库表级 CDC Debezium + Connect
Java 团队的轻中量流处理 Streams
SQL 快速分析/原型 ksqlDB
重状态、事件时间、大规模作业 Flink
复杂图计算/批流一体 Spark/Flink

本章小结

思考题

  1. Connect 的 source offset 存在哪里?为什么这比存在内存里可靠?
  2. Streams 的 changelog topic 删除后会发生什么?如何恢复?
  3. 用 Streams 实现一个“每 1 分钟统计各接口 5xx 次数并告警”的拓扑,写出伪代码。