Kafka Connect 分布式运行时:Worker、内部 Topic 与故障恢复
凌晨的迁移窗口里,值班同学把第二台 Connect worker 加进集群,希望分担十几个 source 与 sink task。新进程启动正常,GET /connectors 也能看到全部连接器;几分钟后却出现三种互相矛盾的现象:有些 task 从旧 worker 消失后迟迟没有在新 worker 出现,有些 connector 显示 RUNNING 但目标端停止增长,还有一个 source 重启后重复发送了一段数据。团队先后重启 worker、删除 connector、重建业务 topic,故障反而扩大,因为真正决定恢复位置的不是进程目录,也不是 REST 返回的一行 RUNNING,而是 worker group、插件集合以及 Kafka 中的配置、位点和状态记录。
Kafka Connect 是 Kafka 提供的连接器运行时。connector 描述“从哪里读、向哪里写以及怎样转换”,task 是实际执行的并行单元,worker 是承载 connector 与 task 的 JVM 进程。distributed 模式把多个 worker 组成一个逻辑集群:配置由 REST 控制面写入 Kafka,任务由 group 协议分配,source offset 和运行状态也写入 Kafka。于是 worker 可以被替换,但前提是新 worker 能读到同一组内部 topic、加载兼容插件并通过相同安全身份访问依赖。
先分清进程存活、任务运行和数据前进
Connect 故障最容易被一个绿色状态误导。worker 的 HTTP 端口可访问,只能证明 Web 服务线程还活着;connector 的状态是 RUNNING,只表示 connector 实例没有报告失败;task 的状态是 RUNNING,也不等于它正在产生有效吞吐。数据库权限被回收、外部 API 一直超时、sink 遇到可重试错误、source 没有新数据时,都可能出现“状态绿但水位不动”。
排查时至少同时观察四个层面:
worker 是否仍属于预期 group.id,最近是否发生 rebalance。connector 和每个 task 分别运行在哪个 worker,是否有 FAILED 及 trace。source offset 或 sink consumer offset 是否前进,最后提交时刻是否更新。
源端水位、Kafka 业务 topic 尾部和目标端结果是否形成连续链路。
Apache Kafka 的 Connect 用户指南 区分 standalone 与 distributed:前者把 source offset 放在本地文件中,后者将配置、offset 和 status 放进 Kafka topic,并支持自动均衡和容错。个人调试一个文件连接器可以用 standalone;持续 CDC、共享数据管道和需要替换 worker 的环境应使用 distributed。把 standalone 的本地 offset 文件复制到多台机器,既不会得到协调,也可能让多个实例从相同位置重复工作。
standalone 不是“单节点 distributed”。它在一个 JVM 中读取 worker properties 和一个或多个 connector properties,没有分布式 REST 配置日志、worker group 再分配以及 config/status topic。下面用官方包自带的 FileStream 连接器跑通最小链路:
# connect-standalone-lab.properties
bootstrap.servers=localhost:9092
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter.schemas.enable=false
offset.storage.file.filename=/tmp/connect-standalone-lab.offsets
offset.flush.interval.ms=10000
plugin.path=/opt/kafka/libs/connect-file-4.3.1.jar# file-source-lab.properties
name=standalone-file-source
connector.class=FileStreamSource
tasks.max=1
file=/tmp/connect-standalone-input.txt
topic=lab.standalone.eventsprintf 'first\n' >/tmp/connect-standalone-input.txt
/opt/kafka/bin/connect-standalone.sh \
connect-standalone-lab.properties file-source-lab.properties另开终端消费 lab.standalone.events,应看到 first;优雅停止后,/tmp/connect-standalone-lab.offsets 应存在且非空,再追加 second 并重启,应只接续新行。反向实验是停止进程、移走 offset 文件后再启动:预期旧行被重放,证明本地文件就是 source 恢复坐标。实验结束先停止进程,再删除合成 topic、输入文件和 offset 文件。迁往 distributed 不能只复制 connector properties:应停写或记录源水位,导出 standalone offset,按连接器支持的 offset 迁移能力写入新集群,最后用重复/缺口对账证明切点;不支持安全导入时,应使用新 topic 受控重放或重新快照。
用发布版和绝对插件路径搭起实验集群
Apache Kafka 下载页当前提供 4.3.1,官方快速入门要求 Java 17+;下载包与官方 apache/kafka:4.3.1 镜像都可作为实验基线。Kafka 按 Apache License 2.0 发布,第三方 connector 及其传递依赖不自动继承这一许可,上线前仍要核对各自许可证、NOTICE、再分发条件和漏洞记录。既有平台应先按 升级与兼容文档 核对 broker、客户端、连接器及序列化格式,不要因为 Connect worker 无状态化外观就跳过兼容测试。
最透明的本机入口是解压 Kafka 二进制包,先启动一个仅供实验的 broker,再写 distributed worker 配置。示例使用官方包中的 FileStream 插件验证运行时;生产 CDC 插件应从供应商发布页下载并核验校验和、签名、许可证和依赖矩阵。
bootstrap.servers=localhost:9092
group.id=connect-lab
config.storage.topic=_connect-lab-configs
offset.storage.topic=_connect-lab-offsets
status.storage.topic=_connect-lab-status
config.storage.replication.factor=1
offset.storage.replication.factor=1
status.storage.replication.factor=1
key.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
plugin.path=/opt/kafka/libs/connect-file-4.3.1.jar
plugin.discovery=hybrid_warn
rest.port=8083
rest.advertised.host.name=localhost
offset.flush.interval.ms=10000
connector.client.config.override.policy=Allowlist
connector.client.config.override.allowlist=client.id实验环境单 broker 只能使用副本因子 1。共享或生产集群应按可容忍的 broker 故障数设置副本,并手工创建内部 topic;worker 配置中的 replication factor 只影响自动创建,不能修正已存在 topic。plugin.path 应使用绝对路径,并让同一 worker group 中每台机器的插件目录、版本与依赖一致。一个插件目录应只容纳该插件及其私有依赖;Connect 为插件创建隔离类加载器,但 Kafka API、SLF4J 等框架类仍有共享边界,隔离并不保证任意依赖组合都兼容。官方 worker 配置参考 显示 plugin.discovery 默认是 hybrid_warn:反射扫描与 ServiceLoader 并用,并对缺少 ServiceLoader 元数据的插件告警。只有在每个插件都通过发现清单验证后,才可切到启动更快的 service_load;否则插件可能不可用。hybrid_fail 可把供应链缺口前移到启动阶段,代价是一个不合规插件会阻止 worker 启动。
Docker Compose 适合本地多 worker 联调。下面补齐单节点 KRaft Kafka,并让两个 worker 挂载同一份只读基础配置和同一个宿主机输入目录。启动命令分别把 advertised host 改成 connect-a 与 connect-b;这两个名称在 Compose 网络内可解析,worker 间 REST 转发不会回到各自的 localhost。
services:
kafka:
image: apache/kafka:4.3.1
hostname: kafka
ports: ["29092:29092"]
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: CONTROLLER://:9093,PLAINTEXT://:9092,HOST://:29092
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,HOST:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"]
interval: 5s
timeout: 5s
retries: 20
connect-a:
image: apache/kafka:4.3.1
command:
- /bin/bash
- -ec
- |
sed -e 's|bootstrap.servers=localhost:9092|bootstrap.servers=kafka:9092|' \
-e 's|rest.advertised.host.name=localhost|rest.advertised.host.name=connect-a|' \
/etc/kafka/connect-base.properties >/tmp/connect-distributed.properties
exec /opt/kafka/bin/connect-distributed.sh /tmp/connect-distributed.properties
ports: ["8083:8083"]
volumes:
- ./connect-distributed.properties:/etc/kafka/connect-base.properties:ro
- ./connect-data:/data
depends_on:
kafka:
condition: service_healthy
connect-b:
image: apache/kafka:4.3.1
command:
- /bin/bash
- -ec
- |
sed -e 's|bootstrap.servers=localhost:9092|bootstrap.servers=kafka:9092|' \
-e 's|rest.advertised.host.name=localhost|rest.advertised.host.name=connect-b|' \
/etc/kafka/connect-base.properties >/tmp/connect-distributed.properties
exec /opt/kafka/bin/connect-distributed.sh /tmp/connect-distributed.properties
ports: ["8084:8083"]
volumes:
- ./connect-distributed.properties:/etc/kafka/connect-base.properties:ro
- ./connect-data:/data
depends_on:
kafka:
condition: service_healthy把前一段 properties 保存为 connect-distributed.properties,Compose 保存为 compose.yaml,再执行:
mkdir -p connect-data
: > connect-data/connect-input.txt
docker compose up -d
curl -s http://localhost:8083/connector-plugins
curl -s http://localhost:8084/connector-plugins两个插件列表都应包含 FileStreamSource。这个实验依赖共享文件,是为了让 task 从任一 worker 接管后仍能读取同一路径;真实生产文件采集不能假定本地盘跨节点共享,应使用共享存储、节点亲和约束或具备远端源语义的 connector。单 broker 示例只能验证机制,不能验证 broker 容错。
共享环境通常把 worker 部署在虚机、容器平台或 Kubernetes。Connect 的扩缩容单位是 worker 进程,不是 connector;容器编排器只负责进程副本与资源,任务分配仍由 Connect group 完成。将 REST 暴露到公网、让弹性伸缩根据 CPU 每分钟抖动、或让不同版本镜像同时长期存在,都会把控制面波动放大成频繁 rebalance。
三类内部 topic 是运行时的持久状态
启动 worker 前手工创建内部 topic,可以避免 broker 默认的分区数、清理策略和副本数悄悄决定恢复能力:
bin/kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic _connect-lab-configs --partitions 1 --replication-factor 1 \
--config cleanup.policy=compact
bin/kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic _connect-lab-offsets --partitions 12 --replication-factor 1 \
--config cleanup.policy=compact
bin/kafka-topics.sh --bootstrap-server localhost:9092 --create \
--topic _connect-lab-status --partitions 5 --replication-factor 1 \
--config cleanup.policy=compactconfig.storage.topic 保存 connector 配置、task 配置和控制记录。它必须只有一个分区,因为全局配置变化需要单一有序日志;错误地创建为多分区不是“多一点并行”,而是破坏一致顺序。offset.storage.topic 保存 source partition 到 source offset 的映射。分区可以较多,以分散大量 connector 的提交压力;日志压缩保留每个逻辑键的最新值,但旧记录不会在写入瞬间消失。status.storage.topic 保存 connector 和 task 的状态、worker 标识和 generation 等信息,它服务于可观测控制面,不是业务正确性的最终证据。
三类 topic 的名称共同定义一个 Connect 集群身份。两个环境若意外复用同一 group.id 和内部 topic,worker 会互相 rebalance,测试配置可能覆盖共享环境配置;若只复用 group.id 却换了内部 topic,又会形成协调关系与状态来源错位。团队模板应让环境名同时进入 group.id、三个 topic 名和安全主体,并在发布门禁中检查唯一性。
不要用普通消费者“整理”内部 topic,不要改 key/value converter,不要把 cleanup.policy 改成纯 delete,也不要把 topic 当备份文件逐条编辑。需要灾备时,应保护 Kafka 集群本身,并记录 worker 配置、插件制品及内部 topic 的恢复策略。恢复副本时,三个 topic 必须来自相互一致的时间线;只恢复 configs 而丢失 offsets,会让 source 无法知道已读位置。
用 REST 建立正向验证闭环
distributed 模式不在启动命令后附 connector 配置文件,连接器通过 REST 创建和管理。先检查插件可见性,再创建一个合成文件 source:
curl -s http://localhost:8083/connector-plugins
curl -s -X PUT http://localhost:8083/connectors/lab-file-source/config \
-H 'Content-Type: application/json' \
-d '{
"connector.class": "FileStreamSource",
"tasks.max": "1",
"file": "/data/connect-input.txt",
"topic": "lab.connect.events",
"key.converter": "org.apache.kafka.connect.storage.StringConverter",
"value.converter": "org.apache.kafka.connect.storage.StringConverter"
}'
curl -s http://localhost:8083/connectors/lab-file-source/status
printf 'order-1001\norder-1002\n' >> connect-data/connect-input.txt
docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka:9092 --topic lab.connect.events \
--from-beginning --max-messages 2这里使用 PUT /connectors/{name}/config,因为同一个声明可重复应用;POST /connectors 在名称已存在时会冲突。预期状态 JSON 中 connector 与 task 0 都为 RUNNING,消费者依次看到两行合成数据。再查询 GET /connectors/lab-file-source/topics 可确认运行时记录的 topic 使用情况。完整端点与状态码应以 Connect REST API 为准。
真正证明恢复能力的步骤是:从 status 的 worker_id 判断 task 在 connect-a:8083 还是 connect-b:8083,停止对应容器,在另一个宿主端口轮询 status,待 task 转移后执行 printf 'order-1003\n' >> connect-data/connect-input.txt。共享挂载让新 worker 仍看到同一文件和字节内容,预期只新增第三条记录。如果重放前两条,先判断 FileStream 插件的位点是否已在停止前提交,而不是把重复归咎于 Kafka。offset.flush.interval.ms 越大,正常路径提交开销越低,但异常退出时可能重放更长窗口;这正是 source 默认“至少一次”需要下游幂等的原因之一。
REST 成功返回并不保证 task 已开始。创建或更新配置后应等待 task 状态稳定、观察日志中的 assignment,并用端到端数据水位验证。自动化发布脚本应保存请求摘要和响应状态码,出现 409 Conflict 时等待 rebalance 完成后重试;不要在冲突期间循环删除再创建 connector。
反向实验:让插件不一致暴露错误机制
复制实验集群后,从 connect-b 的 plugin.path 移除 FileStream JAR,再重启该 worker。GET /connector-plugins 是 worker 本地视图:请求不同 worker 可能返回不同列表。下一次 rebalance 若把 task 分到缺插件的 worker,常见现象是 task 进入 FAILED,trace 中出现找不到 connector 类、转换器或依赖类;也可能 connector 创建阶段就因配置验证失败而拒绝。
这个实验揭示两个机制。第一,插件不存放在 config topic 中,分布式状态只能恢复“要运行什么”,不能恢复“代码在哪里”。第二,REST 负载均衡器若只做健康检查,可能把创建请求发到插件较少的 worker,造成同一地址下的能力不确定。修复顺序应是暂停 connector 或停止变更、统一插件制品、逐台滚动重启、在每个 worker 本地核对 /connector-plugins,最后恢复 connector 并验证 offset 连续。
另一个反例是把 _connect-lab-configs 预建为多个分区。worker 可能在启动或读写配置时报告内部 topic 不符合要求。此时不能靠重启消除结构错误。实验集群可以停掉所有 worker、删除三个内部 topic 后按正确结构重建,再重新创建 connector;共享环境禁止直接这样做,因为删除 offset topic 等价于销毁恢复坐标。应从一致备份恢复或为新集群使用新名称,经过数据对账后切换。
Rebalance、暂停和重启不是同一个动作
当 worker 加入、退出,connector 增删改,task 数变化或协调超时,Connect 会重新分配工作。rebalance 期间,部分 REST 变更可能返回 409;task 先停止并提交可提交的状态,再在新 worker 初始化。source connector 能否无缝续跑,取决于它是否正确实现 offset;sink connector 则通常依赖 Kafka consumer group offset。scheduled.rebalance.max.delay.ms 默认允许离开的 worker 在一段延迟内返回,这能减少短暂滚动升级造成的任务搬迁,却也会让退出 worker 的 task 暂时无人承载;低 RTO 环境应通过故障演练在迁移次数与接管时间之间取值。插件初始化慢、外部系统断连、JVM 停顿和频繁弹性伸缩都会拉长不可用窗口。
pause 阻止 connector 和 task 继续工作,但保留配置与 offset;resume 从既有状态继续。restart 重新初始化 connector 或 task,适合暂态网络、凭证刷新后恢复;带 includeTasks=true&onlyFailed=true 的重启可以减少健康 task 扰动。删除 connector 会移除配置和状态,但 source offset 的生命周期必须按目标版本 REST 能力和连接器文档单独确认,不能把“名称删除”理解为“位点已清空”。Apache Kafka 4.3 的 REST API 提供 connector offset 查看与管理入口,修改前需要先停止 connector,并把导出的 offset、变更审批和回滚坐标一起保存。
tasks.max 只是上限。一个 MySQL binlog source 常因单一有序日志只能创建一个 task;把 tasks.max 从 1 改成 8 不会自动把 binlog 切八份。某些 sink 能按 topic partition 并行,但并行度还受输入分区数、目标端并发写能力和 connector 实现限制。扩 worker 只能为可拆分的 task 提供槽位,不能突破源端串行边界。
容量规划应把 worker 数量与 task 数、每 task CPU/堆内存、批次大小、序列化成本、外部调用延迟一起看。一个大批次 task 发生长 GC,会影响同 JVM 中其他 connector;高风险 connector 可按安全域、插件依赖或资源特征拆到独立 group.id,用不同内部 topic 隔离故障。代价是更多集群、权限和运维对象,不能为了“统一平台”把所有数据库高权凭证放进同一 worker 池。
offset 与 status 要怎样读才不误判
source offset 是连接器定义的结构,不一定是 Kafka 的数字 offset。文件 source 可能是文件路径与字节位置,数据库 CDC 可能是日志文件、位置、GTID 或 LSN。Connect 只负责持久化键值和在 task 启动时交还;位点是否仍可用于源系统、是否和 schema history 匹配,由 connector 负责。
status topic 记录 RUNNING、PAUSED、FAILED 等控制状态。task 抛出不可恢复异常时,REST status 的 trace 是第一份证据;随后要关联 worker 日志、源端错误、内部 topic 可用性和端到端水位。状态记录是异步写入的,短暂旧值不应触发删除 connector。监控应对“连续多次失败”“无 offset 前进且源端有变更”“反复 rebalance”“worker 数低于容错目标”分别告警,而不是只采集一个 connector state。
变更 offset 是高风险运维动作。正确流程是停止 connector,导出当前 offset,确认源端对应日志仍存在,计算目标位置对数据的影响,执行修改,启动后限定时间观察重复、缺口和 schema 解析,再做分块校验。向后移动通常造成重放,向前移动可能永久跳过变更。即使工具允许,也不能把位点前跳当成清除积压的常规手段。
source EOS 只封住 Kafka 内部的一段窗口
默认 source 交付是至少一次:记录可能已经写入业务 topic,但对应 source offset 尚未提交,worker 此时崩溃就会重放。Kafka Connect 的 source exactly-once 使用 Kafka 事务把 source records 与 source offsets 一起提交,并在新 task generation 启动前隔离旧 generation。它只适用于 distributed 模式;standalone 没有这套分布式事务与 fencing 协议。
新建集群可让所有 worker 直接设置 exactly.once.source.support=enabled。已有集群必须先把所有 worker 滚动到 preparing,确认整组稳定后,再全部滚动到 enabled;混用阶段不能把 connector 当作已经获得 EOS。连接器配置再用 exactly.once.support=required 表示“不支持就拒绝启动”,并为 worker principal 授予所需 transactional ID 与内部 topic 权限。开启后,普通 source 的 offset.flush.timeout.ms 不再控制事务提交;事务超时、broker transaction.state.log.* 可用性和授权成为新的恢复依赖。
正向实验应使用声明支持 EOS 的合成 source,消费者设置 isolation.level=read_committed,在持续写入时强杀 task 所在 worker,再按事件唯一键核对已提交输出与 offset。预期旧 generation 被 fencing,消费者看不到 aborted transaction,恢复后水位连续。反向实验把消费者改为 read_uncommitted,或撤销 transactional ID 权限:前者可能暴露已中止记录,后者应让 task 以事务授权错误失败。这两种证据说明 EOS 不是“零重复”开关。
边界必须写进数据契约:Kafka 事务不能回滚源数据库已提交事务,也不能让任意 sink 数据库、对象存储或外部 API 自动加入同一事务。sink connector 的 Kafka consumer offset 与外部写入之间仍取决于该 connector 的幂等或事务实现;多集群复制、非 read_committed 消费者以及事务外副作用也不在 source EOS 内。即便启用 EOS,下游仍应保留稳定 key、业务版本和可回放对账。
REST、Kafka 和外部系统需要三层身份
Connect 至少跨三种安全边界:管理员调用 REST;worker 管理内部 topic 和 group;connector task 访问业务 topic或外部系统。Kafka 官方文档指出,安全集群中的管理客户端、source producer 与 sink consumer 可能需要分别配置身份,worker 级 security.*、producer.*、consumer.* 以及 connector 级 override 不能混为一组。
Kafka 4.3.1 的 connector.client.config.override.policy 默认仍是 All,即 connector 配置可覆盖任意 Kafka client 属性;这对共享 REST 写入口过于宽松。前面的 worker 配置显式改为 Allowlist,并通过 connector.client.config.override.allowlist 只开放 client.id,因此 connector 可提交 producer.override.client.id、consumer.override.client.id 或 admin.override.client.id 来改善审计归属,其他 override 会在配置验证阶段被拒绝。若环境不需要任何 connector 级覆盖,使用 None 更小;Principal 已弃用,不应作为新部署基线。
可执行的反向检查是在现有 connector JSON 中临时加入下面任一未授权字段,再调用配置校验或 PUT:
{
"producer.override.bootstrap.servers": "untrusted.example:9092",
"producer.override.sasl.jaas.config": "<redacted-test-value>"
}预期响应指出 override 不在 allowlist,connector 不应启动。这个实验只验证拒绝路径,不要填写真实凭证,也不要搭建外部接收端。bootstrap.servers 可把 source 数据导向其他集群;sasl.jaas.config、sasl.login.class、interceptor、serializer 等可加载类或携带凭证;acks、重试、超时和隔离级别则可能悄悄改变交付语义。即使某项业务上必须开放,也应按精确属性评审,限制 REST 写权限并隔离高风险 worker 池,而不是允许配置提交者任意重写 producer、consumer 或 admin client。
REST 默认不应直接暴露给不受信网络。前置网关应启用 TLS、身份认证、按 connector 或环境授权、请求体大小限制和审计;GET /connector-plugins 虽不含密码,也会暴露供应链清单。创建与更新配置的响应可能回显连接串或凭证,因此日志、APM body capture 和工单附件必须脱敏。
连接器配置不要提交明文密码。Kafka 支持 ConfigProvider 从外部来源解析配置引用;密钥系统中的读取身份授予 worker,而 connector 配置只保存引用。仍需注意:解析后的 secret 会进入进程内存,插件也在同一 JVM 安全边界内。来源不可信的 connector JAR 获得的不只是一个扩展点,而是 worker 所有可读凭证和网络权限。
Kafka ACL 至少按资源区分:worker group、三个内部 topic、source 写入的业务 topic、sink 读取的业务 topic以及需要时的 transactional ID。开启 source exactly-once 后还要增加事务相关权限,具体模式以 Connect ACL 要求 为准。不要为了省事给 worker 集群管理员权限;分离环境与 connector 池,可以缩小凭证泄露和误写的半径。
从崩溃到可证明恢复
恢复不从“重启一下”开始,而从判断丢失了哪一层状态开始:
| 现象 | 首要证据 | 常见原因 | 恢复动作 |
|---|---|---|---|
| REST 不通 | 进程日志、端口、健康检查 | JVM 退出、端口冲突、错误 advertised 地址 | 恢复 worker;不动 connector 配置 |
connector FAILED | status trace、插件列表 | 配置错误、插件缺失、凭证或网络失败 | 修复依赖后仅重启失败实例 |
| 全集群频繁迁移 | group 日志、worker 数变化 | 自动伸缩抖动、GC、网络或版本不一致 | 稳定成员,减少并发变更 |
| 状态绿但数据不动 | offset、源端水位、业务 topic 尾部 | 源无数据、外部阻塞、静默重试 | 沿数据链逐段定位 |
| 重启后重复 | 重启前后 source offset、提交间隔 | 异常退出前记录未提交 | 下游幂等并缩短可接受重放窗 |
| 重启后无法续跑 | offset 内容、源日志保留 | offset topic 丢失或源日志过期 | 从一致快照重建,禁止盲目前跳 |
worker 机器损坏但 Kafka 与插件制品仍在时,启动一个配置相同的新 worker即可触发任务迁移。Kafka 内部 topic 损坏则是控制面灾难:先冻结变更,验证副本与备份,恢复三个 topic 的一致状态,再恢复同版本插件和 worker。源端日志已经过期时,即便 offset 完好也无法续读,必须由具体 connector 的重新快照或恢复协议接管。
升级采用兼容池验证和滚动替换。先在隔离 group.id 上用生产同款插件与合成数据验证配置解析、序列化、offset 恢复和失败重试;再对正式池逐台替换,观察每轮 rebalance 完成后才继续。回滚不仅是换回镜像,还要确认新版本是否写入旧版本不能理解的配置、offset 或插件私有状态。插件制品、worker 配置和内部 topic 的恢复点共同组成回滚包。
团队把 Connect 当平台时要守住的线
单机 worker 适合个人验证,优点是简单,缺点是没有进程级容错。两台以上 distributed worker 适合共享环境,能自动迁移 task,但必须承担内部 topic、插件一致性、协调和安全管理。按业务域拆多个 Connect 集群能隔离资源和权限,代价是平台对象增多。托管 Connect 降低 worker 运维负担,却不会替团队决定 connector 数据权限、序列化契约、源端日志保留、目标幂等和退出路径。
平台目录应为每个 connector 保存 owner、数据分级、插件与版本、配置仓库位置、源和目标、预期并行度、恢复点类型、最大可接受重放窗口、下游幂等键、告警和下线日期。配置变更走代码评审与自动校验:检查 connector class 白名单、topic 命名、secret 引用、tasks.max 上限、converter 契约和危险 SMT;部署系统通过 REST 应用,而不是允许多人在共享控制台直接改 JSON。
成本不只有 worker CPU。全量快照和重放消耗源端 IO、Kafka 网络与存储;offset/status 提交频率增加内部 topic 写入;错误重试可能打满外部 API;高基数指标和完整 payload 日志会推高观测成本并泄露数据。容量看板至少包括 task 数、worker 资源、rebalance 时长、批次延迟、offset 提交失败、业务 topic 增长、错误重试和源/目标水位差。
连接器退出应先停止新配置变更,确认源端不再需要同步或替代链路已追平,暂停并记录最终 offset,完成数据对账,再删除 connector、回收 ACL 与外部凭证、清理业务 topic 和插件。内部 topic 属于整个 worker group,不能因为下线一个 connector 而删除。若未来可能审计恢复,offset 导出和配置摘要应按数据保留制度归档,但不得保存明文密码。
用故障演练验收而不是用启动日志验收
一个可接管的 Connect 集群,应能重复回答这些问题:任意 worker 退出后 task 在哪里恢复;三类内部 topic 的分区、副本和压缩策略是否符合设计;每台 worker 的插件清单是否一致;offset 停止前后如何证明连续;REST 谁能写、写了什么;源端日志过期时由谁批准重新快照;connector 删除后凭证、topic 和审计记录由谁清理。
最小演练按顺序执行:创建合成 connector,写入两条数据并记录 offset;停止承载 task 的 worker;等待迁移后写第三条;故意撤销源或目标权限并捕获 FAILED 或重试证据;恢复权限并只重启失败 task;最后暂停、导出 offset、删除 connector 和实验业务 topic。预期结果不是“始终零重复”,而是重复窗口可解释、任务归属可查询、故障证据可关联、恢复后水位继续前进,且清理不伤及同集群其他 connector。
删除实验对象可使用:
curl -s -X DELETE http://localhost:8083/connectors/lab-file-source
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--delete --topic lab.connect.events只有销毁整个实验 worker group 且确认没有其他 connector 时,才能停止全部 worker并删除 _connect-lab-configs、_connect-lab-offsets、_connect-lab-status。共享环境的清理必须以 connector 为单位;内部 topic 是集群恢复协议的一部分,不是临时缓存。
