四个交互模块:① Worker/Task 拓扑与故障切换 ② Connector 状态机 ③ Source→Kafka→Sink 端到端动画 ④ SMT 前后对比。
3 Worker 上跑两个 Connector,共 6 个 Task。点击「Kill Worker」观察 Task 自动 Rebalance。
curl -X PUT http://localhost:8083/connectors/<name>/pause
点击「下单」,模拟一条订单从 MySQL → Debezium Source → Kafka → ES Sink → Elasticsearch 的旅程。
类:io.confluent.connect.jdbc.JdbcSourceConnector
模式:bulk / incrementing / timestamp / timestamp+incrementing
痛点:感知不到 DELETE
典型:维度表全量同步、事件流增量同步
类:io.debezium.connector.mysql.MySqlConnector
读:直接订阅 binlog / WAL / oplog
支持:MySQL / PG / Mongo / Oracle / SQL Server
典型:实时数仓、缓存失效、双写一致
类:io.confluent.connect.elasticsearch.ElasticsearchSinkConnector
关键:key.ignore=false 用 PK 幂等覆盖
典型:实时搜索、日志检索
类:io.confluent.connect.s3.S3SinkConnector
格式:JSON / Avro / Parquet
分区:TimeBasedPartitioner(按天 / 小时)
典型:冷数据归档、ODS 层入湖
类:com.mongodb.kafka.connect.MongoSinkConnector
策略:upsert by document key
典型:业务系统二级存储
类:io.confluent.connect.hdfs.HdfsSinkConnector
格式:Avro / Parquet / 文本
典型:Hadoop 生态数据落地
| Converter | 体积 | Schema | 需 Registry | 适合 |
|---|---|---|---|---|
StringConverter | 小 | 无 | ❌ | 纯文本 |
JsonConverterschemas.enable=false | 大 | 无 | ❌ | 调试友好、内部小工具 |
JsonConverterschemas.enable=true | 巨大 | 有(每条带) | ❌ | 不推荐:太占带宽 |
AvroConverter | 极小 | 强 | ✅ | ⭐ 跨团队、强 schema 演进 |
ProtobufConverter | 极小 | 强 | ✅ | gRPC 生态、高性能 |
JsonSchemaConverter | 中 | 强 | ✅ | Web/JS 生态友好 |
单进程跑 connector,配置写本地 properties
Offset 存本地文件
无 HA、改配置要重启进程
适合:开发调试、IoT 网关、教学
bin/connect-standalone.sh \ config/connect-standalone.properties \ config/source-1.properties
多 Worker 组成集群,配置 / Offset / 状态都在 Kafka 内部 Topic
HA:单 Worker 挂自动 Rebalance
REST API 热更新无需重启
适合:所有生产部署
bin/connect-distributed.sh \ config/connect-distributed.properties
选一组 SMT,看 Debezium 原始消息被「翻译」成什么。
| 动作 | 方法 + 路径 |
|---|---|
| 列出已注册 Connector | GET /connectors |
| 列出所有可用 Plugin | GET /connector-plugins |
| 创建 | POST /connectors + JSON body |
| 查看配置 | GET /connectors/{name}/config |
| 查看状态 | GET /connectors/{name}/status |
| 暂停 | PUT /connectors/{name}/pause |
| 恢复 | PUT /connectors/{name}/resume |
| 重启 connector + tasks | POST /connectors/{name}/restart?includeTasks=true |
| 仅重启 FAILED tasks | POST /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 名 | 类型 | 干什么 | 典型用例 |
|---|---|---|---|
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 转扁平表 |
ExtractNewRecordState | — | Debezium 解包:把 after 提出来 | 解 envelope |
| # | 坑 | 表现 | 修复 |
|---|---|---|---|
| 1 | tasks.max 设过大 | 多余 Task 空转 | Sink 端 ≤ 订阅 topic 总分区数 |
| 2 | 手动改 connect-offsets | 数据重复 / 丢失 | 必须先 stop connector,用 REST /offsets |
| 3 | JDBC Source 感知不到 DELETE | 下游 ES 永久残留 | 切 Debezium |
| 4 | Schema 演进破坏 Sink | ES mapping 拒收 | Schema Registry + BACKWARD 策略 |
| 5 | DLQ 自己消费自己 | 死循环 | DLQ 必须独立 Connector / 集群 |
| 6 | plugin.path 路径错 | ClassNotFound | 必须放在子目录,不平铺 |
| 7 | group.id 多集群撞车 | 配置 / Offset 互相覆盖 | group.id 每集群唯一 |
| 8 | 内部 Topic 副本数 1 | Worker 重启丢配置 | 3 个内部 Topic 副本数 ≥ 3 |
原理:Worker 把 Source 产出 + source offset 写入同一个事务,崩溃 → 整个事务 abort → 下次重做。
| 阶段 | exactly.once.source.support | 含义 |
|---|---|---|
| 1 | disabled(默认) | At Least Once:崩溃时一段消息可能重复 |
| 2 | preparing | 过渡:所有 Worker 已升级,但还在用旧路径 |
| 3 | enabled | 启用 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_connect.md · code/start_connect.sh · code/jdbc_source_config.jsoncode/debezium_mysql_config.json · code/elastic_sink_config.json · code/connect_rest_demo.sh