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 条消息

两个结论:

「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 进程还在不在。

生产边界

动手

  1. 用 exp03a_connect.sh 跑一遍 source + sink,strings _connect.offsets 确认能看到 source 的文件位置;再追加几行,看位置有没有推进。
  2. 删掉 offset 文件但保留 topic,重启 worker,观察 sink 文件从几行变成几行,解释多出来的那几行从哪来。
  3. 把 source 的 SMT 从 hoist,addTs 改成只留 addTs,重启,把真实报错抄下来,说明是哪一层拒绝了它。

自测

  1. source 的 offset 和 sink 的 offset 分别是什么含义?为什么同一条链路里它们是两套东西?
  2. standalone 和 distributed 的 offset 分别存在哪?这个差别在故障恢复时意味着什么?
  3. SMT 在链路的哪一段执行?为什么 InsertField 会拒绝一个字符串记录?
  4. sink 目录不存在导致的失败,和一条坏消息导致的失败,Connect 的处理方式有什么区别?
  5. 重启后 sink 文件从 3 行变成 8 行,最可能的原因是什么?该去哪里查?

↓ 下一步:03 章 · 事件时间与水位线 —— 数据进了 topic 之后,怎么在乱序里把窗口算对。

进入 keel 阅读