🔌 第 16 章 · Kafka Connect 可视化

四个交互模块:① Worker/Task 拓扑与故障切换 ② Connector 状态机 ③ Source→Kafka→Sink 端到端动画 ④ SMT 前后对比。

① Connect Cluster 拓扑:Worker × Task

3 Worker 上跑两个 Connector,共 6 个 Task。点击「Kill Worker」观察 Task 自动 Rebalance。

② Connector 生命周期状态机

UNASSIGNED
RUNNING
PAUSED
FAILED
DESTROYED
curl -X PUT http://localhost:8083/connectors/<name>/pause

③ Source → Kafka → Sink 端到端动画

点击「下单」,模拟一条订单从 MySQL → Debezium Source → Kafka → ES Sink → Elasticsearch 的旅程。

🗄️
MySQL
shop.orders
🔌
Debezium
Source
📨
Kafka Topic
dbz.shop.orders
🔌
ES Sink
🔍
Elasticsearch
orders

⑤ 常见 Connector 一览

🗄 JDBC Source

类:io.confluent.connect.jdbc.JdbcSourceConnector

模式:bulk / incrementing / timestamp / timestamp+incrementing

痛点:感知不到 DELETE

典型:维度表全量同步、事件流增量同步

🐝 Debezium CDC

类:io.debezium.connector.mysql.MySqlConnector

读:直接订阅 binlog / WAL / oplog

支持:MySQL / PG / Mongo / Oracle / SQL Server

典型:实时数仓、缓存失效、双写一致

🔍 Elasticsearch Sink

类:io.confluent.connect.elasticsearch.ElasticsearchSinkConnector

关键:key.ignore=false 用 PK 幂等覆盖

典型:实时搜索、日志检索

☁ S3 Sink

类:io.confluent.connect.s3.S3SinkConnector

格式:JSON / Avro / Parquet

分区:TimeBasedPartitioner(按天 / 小时)

典型:冷数据归档、ODS 层入湖

🍃 MongoDB Sink

类:com.mongodb.kafka.connect.MongoSinkConnector

策略:upsert by document key

典型:业务系统二级存储

🧱 HDFS Sink

类:io.confluent.connect.hdfs.HdfsSinkConnector

格式:Avro / Parquet / 文本

典型:Hadoop 生态数据落地

⑥ Converter 选型

Converter体积Schema需 Registry适合
StringConverter纯文本
JsonConverter
schemas.enable=false
调试友好、内部小工具
JsonConverter
schemas.enable=true
巨大有(每条带)不推荐:太占带宽
AvroConverter极小⭐ 跨团队、强 schema 演进
ProtobufConverter极小gRPC 生态、高性能
JsonSchemaConverterWeb/JS 生态友好

⑦ Standalone vs Distributed

🥚 Standalone

单进程跑 connector,配置写本地 properties

Offset 存本地文件

无 HA、改配置要重启进程

适合:开发调试、IoT 网关、教学

bin/connect-standalone.sh \
  config/connect-standalone.properties \
  config/source-1.properties

🏭 Distributed(生产推荐)

多 Worker 组成集群,配置 / Offset / 状态都在 Kafka 内部 Topic

HA:单 Worker 挂自动 Rebalance

REST API 热更新无需重启

适合:所有生产部署

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

④ SMT 变换前后对比

选一组 SMT,看 Debezium 原始消息被「翻译」成什么。

输入(Debezium 原始消息)


    

输出(SMT 后写入下游 topic)


    

⑧ REST API 速查(默认端口 8083)

动作方法 + 路径
列出已注册 ConnectorGET /connectors
列出所有可用 PluginGET /connector-plugins
创建POST /connectors + JSON body
查看配置GET /connectors/{name}/config
查看状态GET /connectors/{name}/status
暂停PUT /connectors/{name}/pause
恢复PUT /connectors/{name}/resume
重启 connector + tasksPOST /connectors/{name}/restart?includeTasks=true
仅重启 FAILED tasksPOST /connectors/{name}/restart?onlyFailed=true&includeTasks=true
修改配置(PUT,仅 config 部分)PUT /connectors/{name}/config
删除DELETE /connectors/{name}
查看 Source Offset(3.6+)GET /connectors/{name}/offsets
重置 Offset(3.6+)DELETE /connectors/{name}/offsets
查看 / 改 logger 等级GET/PUT /admin/loggers/{package}

💡 完整命令集合见 16_connect/code/connect_rest_demo.sh

⑩ 常用 SMT 速查(Single Message Transforms)

SMT 名类型干什么典型用例
InsertField$Key / $Value插入字段(topic / partition / offset / timestamp)加审计字段
ReplaceField$Key / $Value改字段名 / 删字段适配下游 schema
MaskField$Value把字段值置 null/0脱敏(手机号、身份证)
ValueToKey把 value 字段提到 key让 Sink 用业务主键去重
ExtractField$Key / $Value从结构里抽一个子字段当整条 value提取 Debezium 的 after
Cast$Value类型强转 int → string 等schema 兼容
TimestampRouter按时间戳路由 topic 名按天分 topic
RegexRouter按正则改 topic 名去前缀重命名
Filter满足条件的丢弃只保留 op=u(update)
Flatten$Value把嵌套结构拍平嵌套 JSON 转扁平表
ExtractNewRecordStateDebezium 解包:把 after 提出来解 envelope

⑪ 生产踩坑速查(命中即心碎)

#表现修复
1tasks.max 设过大多余 Task 空转Sink 端 ≤ 订阅 topic 总分区数
2手动改 connect-offsets数据重复 / 丢失必须先 stop connector,用 REST /offsets
3JDBC Source 感知不到 DELETE下游 ES 永久残留切 Debezium
4Schema 演进破坏 SinkES mapping 拒收Schema Registry + BACKWARD 策略
5DLQ 自己消费自己死循环DLQ 必须独立 Connector / 集群
6plugin.path 路径错ClassNotFound必须放在子目录,不平铺
7group.id 多集群撞车配置 / Offset 互相覆盖group.id 每集群唯一
8内部 Topic 副本数 1Worker 重启丢配置3 个内部 Topic 副本数 ≥ 3

⑫ Exactly Once Source(Connect 3.3+)

原理:Worker 把 Source 产出 + source offset 写入同一个事务,崩溃 → 整个事务 abort → 下次重做。

阶段exactly.once.source.support含义
1disabled(默认)At Least Once:崩溃时一段消息可能重复
2preparing过渡:所有 Worker 已升级,但还在用旧路径
3enabled启用 EOS Source:消费者必须 read_committed 才看不到 abort 中间值

⚠ 切换流程必须 disabled → preparing → enabled,逐步生效。Sink 端 EOS 仍依赖下游幂等(如 ES 用 docId 覆盖、JDBC 用 PK upsert)。

⑨ 与其他集成方案对比

方案优点缺点典型场景
Kafka Connect生态丰富、运维统一、HA 自带JVM 依赖、内存大Kafka 生态全栈
Flink CDC流处理 + CDC 一体、EOS需 Flink 集群、Sink 不如 Connect实时数仓
Apache NiFi拖拽 UI、灵活度极高学习陡、社区小企业 ETL
Logstash / Fluentd轻量、日志生态强不适合事务性数据日志收集
Airbyte现代 UI、SaaS 模型实时性弱数据仓库 ELT
Debezium 直接当 lib无 Connect 框架开销自己写部署 / HA嵌入式 CDC
📚 第 16 章配套:16_connect.md · code/start_connect.sh · code/jdbc_source_config.json
code/debezium_mysql_config.json · code/elastic_sink_config.json · code/connect_rest_demo.sh