KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
02 · 管道拓扑与 Kafka Connect — keel 龙骨
第 01 章的解码结果还只是一个 SELECT 的返回值。这一章讲怎么把它变成一个跑了就不停、断了能续的连接器,以及「续」这件事的状态到底写在哪。
第 01 章的解码结果还只是一个 SELECT 的返回值。这一章讲怎么把它变成一个跑了就不停、断了能续的连接器,以及「续」这件事的状态到底写在哪。
现场
一个需求:把服务器上某个日志文件的变化实时送进 Kafka,再原样落回另一个目录,中间要打上一个处理时间戳。
手写的做法是写个脚本 tail -f 然后 producer.send()。能跑。问题出在重启:脚本被 OOM kill 之后,它从文件的哪个位置继续读?如果从头读,整个文件会被重发一遍;如果从末尾读,中间那段丢了。这个「从哪继续」的状态,手写脚本时必须自己落盘——而 Kafka Connect 的整个设计就是围绕这件事展开的。
全链路
flowchart TD
subgraph WORKER["Connect worker(standalone,单进程)"]
direction TB
subgraph SRCT["source task (fs-source-01)"]
F["source-01.txt<br/>原始行"] --> SRC["FileStreamSource<br/>按文件位置读"]
SRC --> SMT["SMT 链:hoist → InsertField<br/>把字符串包成 struct 再加时间戳"]
end
subgraph SNKT["sink task (fs-sink-01)"]
CON["KafkaConsumer<br/>按已提交位点读"] --> SINKF["FileStreamSink<br/>追加写 sink-01.txt"]
end
OS["offset 存储<br/>_connect.offsets<br/>{source 文件位置} / {sink 消费位点}"]
SRCT -. "每个 source 记录记一次文件位置" .-> OS
SNKT -. "每隔 offset.flush.interval.ms 提交一次" .-> OS
end
SMT -->|"JSON 信封"| T["Kafka 集群B topic<br/>connect-file-04"]
T --> CON
SINKF -. "目标目录不存在:task 直接失败" .-> FAIL["ConnectException<br/>Task is being killed"]
style F fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style SRC fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style SMT fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style CON fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style SINKF fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style OS fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style T fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style FAIL fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
精确定义
connector / task:connector 是配置和生命周期的管理者,真正干活的是 task。tasks.max 决定最多拆几个 task 并行——source 侧的并行度通常受「数据源能不能被切开」限制(一个文件只能一个 task)。
source / sink:source 从外部系统读进来写进 topic,sink 从 topic 读出去写到外部系统。两者的 offset 含义完全不同:source 的 offset 是「我在外部系统读到哪了」(本机是文件字节位置),sink 的 offset 是「我在 topic 里消费到哪了」(Kafka 分区位点)。
offset storage:standalone 模式存在本地文件(offset.storage.file.filename),distributed 模式存在内部 topic connect-offsets。这是「断了怎么续」的答案载体。
converter:决定 topic 里字节和 Connect 内部记录(Struct/Map)互转的格式。worker 级配置,所有 connector 共享。
SMT(Single Message Transform):作用在记录被 converter 序列化之前、在 source task 产出之后。链式执行,前一个的输出是后一个的输入。
一次完整运行
原始输出见 lab/evidence/stream-data-pipeline/03-connect.txt。源文件追加 4 行,source 配了 hoist + InsertField 两个 SMT,从 topic 读回来是这样:
{"raw_line":"o-1001,customer=1,amount=10.00","ingested_at":1791189951345}
{"raw_line":"o-1002,customer=1,amount=20.00","ingested_at":1791189951345}
{"raw_line":"o-1003,customer=2,amount=30.00","ingested_at":1791189951345}
{"raw_line":"o-1004,customer=3,amount=40.00","ingested_at":1791189951345}
原始的裸字符串被 hoist 包成了 raw_line 字段,InsertField 又加上 ingested_at。sink 侧落地的文件内容与之一致:
{ingested_at=1791189951345, raw_line=o-1001,customer=1,amount=10.00}
offset 到底存在哪
worker 的 offset 文件不是纯文本,strings 能看出结构:
java.util.HashMap
w["fs-source-01",{"filename":".../source-01.txt"}]uq
{"position":124}x
sink-01.txt 有 4 行、每行 31 字节,正好 124。也就是说 source 的 offset 是「文件读到第几个字节」,不是 Kafka 位点——因为 source 根本不消费 topic,它读的是文件系统。第一次运行后整个 offset 文件里只有 source 一条记录,sink 的那条要等它按 offset.flush.interval.ms 提交后才会出现。
这个细节很容易被忽略:同一条链路里,两端的 offset 语义是两套。排查「重启后重复」时,先确认是哪一端的位点没落地。
断了怎么续
这是本章的核心实验,原始输出见 03-connect-resume.txt。第一次运行末尾追加 3 行,停掉 worker;保留 offset 文件重启,再追加 2 行:
第一次运行后:sink 3 行,offset 文件 position=124(对应前面那轮)
第二次运行后:sink 5 行
{raw_line=r-1} {raw_line=r-2} {raw_line=r-3} {raw_line=r-4} {raw_line=r-5}
offset 文件:position 12 → 20
sink 这次从哪读:Setting offset for partition connect-file-04-0 to the committed offset {offset=3}
topic 里一共 5 条消息
两个结论:
- source 没重发。位置从 12 推进到 20,说明它接着上次的字节位置读,前 3 行没有被重新投递。
- sink 没重读。它从已提交的
offset=3续读,只处理了r-4、r-5,最终文件恰好 5 行而不是 8 行。
「5 行」这个数字就是断点续传成立了:没重、没漏。如果 offset 存储被清空,sink 会因为 auto.offset.reset=earliest 把整个 topic 重读一遍,文件变成 8 行——重复不是 Connect 的 bug,是位点丢了。
SMT 在哪一段生效,以及它的类型约束
SMT 作用在 source task 产出记录之后。这个位置带来一个容易踩的约束:InsertField / MaskField 这类「往结构里加字段」的变换,要求记录值是 Struct,而 FileStreamSource 产出的是裸字符串。
不先包一层直接上 InsertField,task 会当场炸:
ERROR WorkerSourceTask{id=fs-source-01-0} Task threw an uncaught and unrecoverable exception.
org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
Caused by: org.apache.kafka.connect.errors.DataException:
Only Struct objects supported for [field insertion], found: java.lang.String
所以链路里加了 HoistField 先把字符串包成 {raw_line: "..."},InsertField 才有地方插字段。Tolerance exceeded in error handler 这行也值得记住:SMT 里抛的错走的是容错处理器(errors.tolerance / errors.retry.timeout),重试耗尽后 task 被杀——它不是「跳过这条坏消息」,是整条 task 停摆。
同样的道理,想在管道里给身份证号打码,得先让记录是结构化的(比如上游发 JSON),MaskField 才有字段可打;对着裸字符串配 MaskField 会得到和上面同类的报错。
失败注入:目标目录不存在
把 sink 的 file 指向一个不存在的目录,重启 worker,原始输出见 03-connect-failure.txt:
ERROR [fs-sink-bad|task-0] WorkerSinkTask{id=fs-sink-bad-0}
Task threw an uncaught and unrecoverable exception.
Task is being killed and will not recover until manually restarted
org.apache.kafka.connect.errors.ConnectException:
Couldn't find or create file '.../no_such_dir/sink-bad.txt' for FileStreamSinkTask
at org.apache.kafka.connect.file.FileStreamSinkTask.start(FileStreamSinkTask.java:74)
Caused by: java.nio.file.NoSuchFileException: ...\no_such_dir\sink-bad.txt
at org.apache.kafka.connect.file.FileStreamSinkTask.start(FileStreamSinkTask.java:70)
日志里 Creating task fs-sink-bad-0 在整个观察窗口内只出现 1 次。失败发生在 task.start(),worker 的处置是「杀掉 task,等人工重启」,并没有像处理记录级错误那样自动重试。这和上面的 SMT 失败不同:SMT 失败走容错处理器(会重试、计入 tolerance),而 task 启动就抛的异常没有容错路径。
运维含义很直接:这类失败不会自愈,必须有人发现并修。所以 task 状态(RUNNING / FAILED / UNASSIGNED)必须进监控,不能只看 worker 进程还在不在。
生产边界
- 教学替身 vs 真实依赖:本机是 standalone(位点在本地文件、单进程)。生产用 distributed 模式时,位点存在内部 topic
connect-offsets,connector 配置存在connect-configs,task 可以分布到多台 worker 上,失败会重新分配。standalone 的失败就是全停。 - 要盯的指标:每个 connector 的 task 数与状态、source 的
source-record-poll-rate、sink 的sink-record-active-count与消费 lag、以及 offset 提交是否在推进。 - 失败策略:记录级坏数据用
errors.tolerance=all+ 死信队列(把解析失败的原消息写到另一个 topic),别让它杀掉 task;task 级失败(目录、权限、字段类型)要告警到人。 - 环境坑:worker 的
plugin.path指向一个含上百个 jar 的目录时,本机实测启动 45 秒都没打出第一行日志;去掉plugin.path、让连接器走-cp "libs/*"的 classpath 后秒起。生产环境按官方方式把连接器单独放一个目录即可,别指到整套发行包里。
动手
- 用
exp03a_connect.sh跑一遍 source + sink,strings _connect.offsets确认能看到 source 的文件位置;再追加几行,看位置有没有推进。 - 删掉 offset 文件但保留 topic,重启 worker,观察 sink 文件从几行变成几行,解释多出来的那几行从哪来。
- 把 source 的 SMT 从
hoist,addTs改成只留addTs,重启,把真实报错抄下来,说明是哪一层拒绝了它。
自测
- source 的 offset 和 sink 的 offset 分别是什么含义?为什么同一条链路里它们是两套东西?
- standalone 和 distributed 的 offset 分别存在哪?这个差别在故障恢复时意味着什么?
- SMT 在链路的哪一段执行?为什么
InsertField会拒绝一个字符串记录? - sink 目录不存在导致的失败,和一条坏消息导致的失败,Connect 的处理方式有什么区别?
- 重启后 sink 文件从 3 行变成 8 行,最可能的原因是什么?该去哪里查?
↓ 下一步:03 章 · 事件时间与水位线 —— 数据进了 topic 之后,怎么在乱序里把窗口算对。