跳到主要内容

Kafka 消息格式

FZS 支持将 Kafka 作为备端(目标端类型 kafka),把源端的全量数据与增量变更以 JSON 消息实时写入 Kafka Topic。

本文面向从 Kafka 中取数的对接方,说明消息的整体约定、字段含义与取值规则,便于下游解析与落库。

提示

如果你只需要「拿到消息就能用」,建议直接看:

  1. 消息约定 —— Topic 怎么生成、消息写在第几个分区、有没有 Key
  2. 字段值编码规则 —— null / 十六进制 / 字符集转换
  3. Debezium JSON 格式 —— 默认消息格式与各字段含义

适用场景​

场景说明
全量同步源端存量数据逐表读出,按行写入 Kafka(可多行合并为一条消息)
增量同步源端的 insert / update / delete / DDL 逐条写入 Kafka
同构 / 异构支持 Oracle、MySQL、SQL Server、PostgreSQL、达梦、GaussDB、OceanBase、Sundb 等源端
备注

FZS 的写入方向是「源端数据库 → FZS → Kafka」。若 Kafka 是源端(FZS 从 Kafka 消费 Debezium 消息),则属于另一种链路类型,不在本文范围内。

备端配置项​

在 FZS Web 创建链路时选择 Kafka 作为备端,或直接通过 Agent 配置文件下发:

Web 界面配置项配置文件 Key默认值说明
Kafka 消息格式format_jsonDEBEZIUM_JSONDEBEZIUM_JSON 为 Debezium 兼容格式;KZS_JSON 为 FZS 默认格式
Kafka Topickafka_topic空填写后所有消息统一写入该 Topic;不填则按 用户名 + 分隔符 + 表名 自动生成
Kafka Topic 分隔符kafka_topic_delimiter.Topic 自动生成时使用的分隔符
目标端连接串tgt_login—取 @ 之后的地址作为 bootstrap.servers,例如 user/pass@kafka1:9092,kafka2:9092
额外参数kafka.<参数名>—以 kafka. 为前缀的额外参数会透传给 librdkafka,例如 kafka.security.protocol=SASL_PLAINTEXT
注意

消息格式一旦确定不要中途切换。DEBEZIUM_JSON 与 KZS_JSON 的报文结构完全不同,切换后同一条消费链路会遇到两种结构混存。

Topic 命名规则​

  1. 配置了「Kafka Topic」:所有库表的变更都写入该统一 Topic。
  2. 未配置:每条消息按其所属表写入独立 Topic,名字为 用户名(owner) + 分隔符 + 表名,例如 HR.EMPLOYEES。
  3. 若该对象在链路中配置了目标端对象映射(用户/表映射),则直接使用映射后的目标表名作为 Topic 名(不再带 owner 前缀)。
注意

FZS 只负责写入,不会创建 Topic。请提前创建好所需 Topic,或在 Broker 上开启 auto.create.topics.enable。

消息约定​

对接方在编码消费逻辑前,请先确认以下几点:

项目约定说明
消息粒度一条消息 = 一次变更(或一批同表多行插入)多行插入时:KZS_JSON 合并为一条数组消息,DEBEZIUM_JSON 拆成多条消息逐行发送
分区固定写入分区 0FZS 显式指定 partition=0,即使 Topic 有多个分区也只有分区 0 有数据
顺序单分区内严格有序与源端事务提交顺序一致
消息 Key仅 DEBEZIUM_JSON 格式有KZS_JSON 格式不写 Key;分区也不按 Key 计算
消息体编码UTF-8 紧凑 JSON无多余空格/换行
压缩lz4生产端默认 compression.type=lz4
单条上限128 MB生产端 message.max.bytes=134217728;消息体超过 130,000,000 字节会被跳过并打错误日志
发送失败重试3 次仍失败的消息会被丢弃并打错误日志,请关注 FZS 日志
刷盘每个数据文件处理完 flush 一次保证不丢消息
备注

由于分区固定为 0,Topic 分区数建议设置为 1。增加分区不会提高消费并行度(其它分区永远没有数据)。

不发送的消息​

以下变化不会写入 Kafka,对接方无法从 Kafka 中感知:

  • 事务开始 / 事务结束消息(不发送事务边界标记,消息中也没有事务号)
  • OP_DML_ORP、OP_DML_EXP、OP_CTL_TABE 等内部流程消息
  • 未纳入同步范围的对象
注意

无法判断事务边界:同一事务的多条变更会连续写入,但没有「事务开始/结束」标记,消息里也不带 XID。若下游需要事务级一致性,请自行按 SCN(KZS_JSON)或 ts_ms(DEBEZIUM_JSON)分组处理,或让 FZS 直接写目标库。

通用字段​

下表汇总两类格式中出现的字段含义(ts_ms、snapshot 为 DEBEZIUM_JSON 专有;SCN/SCNTIME、OWNER、TABLE、DSQL 等为 KZS_JSON 专有):

字段类型含义
SCN / SCNTIMEstring源端系统变更号(decimal 字符串)。值越大表示变更越晚,可用于排序与延迟判断
OP / opstring操作类型,见下表
HS_RIDstringFZS 行标识,18 位 Base64(字母表 A-Za-z0-9+/),由「数据对象ID(6) + 文件号(3) + 数据块号(6) + 槽号(3)」编码而成。可理解为源端 ROWID 的等价物,同一行的变更通常保持一致,可用于定位与关联
srcIdnumber目标端(链路)在 FZS 内部的编号 tgt_id,多目标写入同一 Kafka 时用于区分来源
OWNER / schemastring源端用户名 / Schema
TABLE / tablestring源端表名
ts_msnumber变更时间戳(毫秒)。注意:由事务提交时间(秒级)换算而来,末 3 位恒为 000
snapshotstring是否全量(快照)阶段的数据:"true" / "false",是字符串不是布尔值

操作类型取值:

值含义对应源端操作
I / c插入INSERT(含多行插入 OP_DML_QMI)
U / u更新UPDATE
D / d删除DELETE
LDDL建表 / 改表 / 删表等(仅 KZS_JSON 格式使用 OP="L")
B / E事务开始 / 结束不会发送,仅保留定义

格式一:Debezium JSON(默认)​

与 Debezium 的消息信封(envelope)结构保持一致,方便已有 Debezium 消费程序复用。

顶层结构​

{
"before": { "...": "变更前镜像,无则为 null" },
"after": { "...": "变更后镜像,无则为 null" },
"source": { "...": "来源信息" },
"op": "c | u | d",
"ts_ms": 1735689600000
}
字段类型含义
beforeobject | null变更前镜像。插入时为 null;删除时为被删除行的旧值;更新时为旧值WHERE列
afterobject | null变更后镜像。插入时为整行新值;删除时为 null;更新时为新值DATA列
sourceobject来源描述,见下表
opstringc 插入 / u 更新 / d 删除
ts_msnumber同 source.ts_ms,毫秒时间戳

source 字段​

字段类型含义
ts_msnumber事务提交时间换算的毫秒时间戳(秒级精度)
snapshotstring"true" 表示全量(快照)阶段,"false" 表示增量。字符串类型
schemastring源端用户名 / Schema
tablestring源端表名
提示

source.schema + source.table 即源端表。若链路配置了用户/表映射,这里的值仍是源端名称,映射关系请以 FZS 链路配置为准。

消息 Key 规则​

DEBEZIUM_JSON 格式会给每条消息写入 Key(JSON 字符串),内容随操作类型变化:

操作Key 内容
插入 c源端类型为 oracle 时包含 ROWID;表存在主键或唯一索引时,追加主键/唯一键列的值
更新 u源端类型为 oracle 时包含 ROWID;表存在主键或唯一索引时,追加**before(旧值)中的全部列**
删除 d源端类型为 oracle 时包含 ROWID;表存在主键或唯一索引时,追加被删除行的全部列
DDL不写 Key
多行插入 QMI消息值是 Key 数组,逐行对应(第 i 个 Key 对应第 i 行)
注意

Key 中的 ROWID 仅在源端类型为 oracle 时才会出现(ob-oracle、达梦等其它源端即使有 rowid 也不会写入 Key)。因此下游不要假设 Key 一定存在或一定包含主键,建议以 HS_RID(KZS_JSON)/ 表主键自行做幂等。

示例​

插入(INSERT) —— Key:

{"ROWID":"AAAC9gAABkEAAAABXm","EMPLOYEE_ID":"100"}

Value:

{"before":null,"after":{"EMPLOYEE_ID":"100","NAME":"张三","SALARY":"1000.00"},"source":{"ts_ms":1735689600000,"snapshot":"false","schema":"HR","table":"EMPLOYEES"},"op":"c","ts_ms":1735689600000}

更新(UPDATE) —— before 为旧值列,after 仅包含被更新的列:

{"before":{"EMPLOYEE_ID":"100","SALARY":"1000.00"},"after":{"SALARY":"2000.00"},"source":{"ts_ms":1735689601000,"snapshot":"false","schema":"HR","table":"EMPLOYEES"},"op":"u","ts_ms":1735689601000}

删除(DELETE) —— after 为 null:

{"before":{"EMPLOYEE_ID":"100","NAME":"张三","SALARY":"1000.00"},"after":null,"source":{"ts_ms":1735689602000,"snapshot":"false","schema":"HR","table":"EMPLOYEES"},"op":"d","ts_ms":1735689602000}

DDL —— 结构与 DML 不同,无 before / after / op,改由 ddl + tableChanges 描述:

{"source":{"ts_ms":1735689603000,"snapshot":"false","schema":"HR","table":"EMPLOYEES"},"ts_ms":1735689603000,"ddl":"ALTER TABLE HR.EMPLOYEES ADD (EMAIL VARCHAR2(100))","tableChanges":[{"type":"ALTER","id":"\"HR\".\"EMPLOYEES\""}]}
DDL 字段类型含义
ddlstringDDL 语句文本(已把换行替换为空格,语句分隔符替换为 ;)
tableChangesarray对象变更列表,type 为 CREATE / ALTER / DROP,id 为 "用户"."表名"
备注

全量同步阶段的建表 DDL 会被 FZS 转换为目标端可执行的 CREATE TABLE 语句后写入。若源端本身就是 Kafka(Debezium),DDL 消息会原样透传:此时消息体为 <key JSON>\x05<value JSON>,FZS 拆分后分别作为消息 Key 与 Value 写入。

与标准 Debezium 的差异​

本格式是 Debezium 的「简化信封」,字段名一致但内容更少,不能直接喂给标准 Debezium 的 schema 解析器:

项目标准 DebeziumFZS
schema(消息 schema 描述)有(含字段类型定义)没有,需下游自行获取表结构
source 内字段含 connector、db、version、pos、server_id 等只有 ts_ms、snapshot、schema、table
snapshot 类型布尔值字符串 "true" / "false"
快照记录 oprc
ts_ms 精度毫秒秒级(末 3 位为 000)
transaction 字段有(事务信息)没有
op 额外取值r(快照读)、t(truncate)只有 c / u / d
注意

ts_ms 取自交易文件的提交时间;若文件不带事务头(例如全量阶段的文件),会沿用上一个可用值,链路刚启动时可能为 0。请勿用 ts_ms=0 判断真实时间,避免下游按时间分区时全部落到 1970 年。

格式二:FZS 默认格式(KZS_JSON)​

FZS 原生格式,字段更贴近源端概念。消息值始终是一个 JSON 数组,元素个数 = 本消息包含的行数(普通变更 1 个元素,多行插入 N 个元素)。

元素结构​

[{
"SCN": "3760546234",
"OP": "I",
"HS_RID": "AAAC9gAABkEAAAABXm",
"message": { "...": "行的详细内容" }
}]
字段类型含义
SCNstring源端系统变更号
OPstring本消息的操作类型:I / U / D / L
HS_RIDstring行标识(DDL 消息无此字段)
messageobject行的实际内容,见下

message 内字段​

字段出现于含义
OWNER / TABLE全部源端用户名 / 表名
OPU / 多行插入 / L操作类型;注意单行插入 I 与删除 D 的 message 内不含 OP(此时以最外层 OP 为准)
HS_RID全部 DML行标识
SCNI / U / D源端系统变更号
SCNTIME多行插入源端系统变更号(多行插入用的是 SCNTIME,不是 SCN)
srcId全部目标端编号
各列名I / D / 多行插入列值直接以「列名: 值」的键值对平铺在 message 中
WHEREU更新前的旧值(用于定位被更新的行)
DATAU更新后的新值(仅含被修改的列)
DSQLLDDL 语句文本

示例​

插入(INSERT):

[{"SCN":"3760546234","OP":"I","HS_RID":"AAAC9gAABkEAAAABXm","message":{"OWNER":"HR","TABLE":"EMPLOYEES","HS_RID":"AAAC9gAABkEAAAABXm","SCN":"3760546234","srcId":1,"EMPLOYEE_ID":"100","NAME":"张三","SALARY":"1000.00"}}]

更新(UPDATE):

[{"SCN":"3760546240","OP":"U","HS_RID":"AAAC9gAABkEAAAABXn","message":{"OWNER":"HR","TABLE":"EMPLOYEES","OP":"U","HS_RID":"AAAC9gAABkEAAAABXn","SCN":"3760546240","srcId":1,"WHERE":{"EMPLOYEE_ID":"100","SALARY":"1000.00"},"DATA":{"SALARY":"2000.00"}}}]

删除(DELETE):

[{"SCN":"3760546250","OP":"D","HS_RID":"AAAC9gAABkEAAAABXo","message":{"OWNER":"HR","TABLE":"EMPLOYEES","HS_RID":"AAAC9gAABkEAAAABXo","SCN":"3760546250","srcId":1,"EMPLOYEE_ID":"100","NAME":"张三","SALARY":"1000.00"}}]

多行插入(一条消息含多行):

[{"SCN":"3760546200","OP":"I","HS_RID":"AAAC9gAABkEAAAABXm","message":{"OWNER":"HR","TABLE":"EMPLOYEES","OP":"I","HS_RID":"AAAC9gAABkEAAAABXm","srcId":1,"SCNTIME":"3760546200","EMPLOYEE_ID":"100","NAME":"张三"}},
{"SCN":"3760546200","OP":"I","HS_RID":"AAAC9gAABkEAAAABXn","message":{"OWNER":"HR","TABLE":"EMPLOYEES","OP":"I","HS_RID":"AAAC9gAABkEAAAABXn","srcId":1,"SCNTIME":"3760546200","EMPLOYEE_ID":"101","NAME":"李四"}}]

DDL:

[{"SCN":"3760546260","OP":"L","message":{"OWNER":"HR","TABLE":"EMPLOYEES","OP":"L","DSQL":"ALTER TABLE HR.EMPLOYEES ADD (EMAIL VARCHAR2(100))","srcId":1}}]
备注

KZS_JSON 格式没有 snapshot 之类字段,无法区分全量还是增量。如果需要区分,请使用 DEBEZIUM_JSON 格式(看 source.snapshot)。

字段值编码规则​

两种格式共用同一套取值规则:

源端情况输出
列为 NULLnull
列值为空字符串(长度 0)null —— 空串与 NULL 无法区分
文本列(CHAR / VARCHAR / LONG / CLOB)JSON 字符串,按 load.param_charset_converter 配置从源端字符集转成目标字符集
UTF-16 文本列(NVARCHAR2 / NCHAR 等)先按 UTF-16BE 解码,再转成目标字符集
二进制列(RAW / LONG RAW / BLOB 等)小写十六进制字符串(如 48656c6c6f),消费方需自行 hex 解码
数值 / 日期 / 时间戳等非文本列源端提供的文本表示,原样透传(Oracle 通常为 YYYY-MM-DD HH24:MI:SS 形式)
数字型字段(srcId)与时间戳(ts_ms)JSON number,其余字段一律为 JSON string 或 null
提示

所有列值都是字符串或 null,不会输出 JSON 数字/布尔。下游请按列类型自行转换,不要依赖 JSON 自动类型推断。

消费示例​

命令行快速查看(带 Key):

kafka-console-consumer \
--bootstrap-server kafka1:9092 \
--topic HR.EMPLOYEES \
--from-beginning \
--property print.key=true \
--property key.separator=' | '

Python 解析示例:

import json
from kafka import KafkaConsumer

consumer = KafkaConsumer(
"HR.EMPLOYEES",
bootstrap_servers="kafka1:9092",
auto_offset_reset="earliest",
)

for msg in consumer:
key = json.loads(msg.key) if msg.key else None # 仅 DEBEZIUM_JSON 有
value = json.loads(msg.value.decode("utf-8")) # KZS_JSON 是数组

if isinstance(value, list): # KZS_JSON
for row in value:
op, fields = row["OP"], row["message"]
print(row["SCN"], op, fields["OWNER"], fields["TABLE"], fields)
else: # DEBEZIUM_JSON
print(value["op"], value["source"]["schema"], value["source"]["table"],
value["before"], value["after"])

对接注意事项​

  1. 提前创建 Topic:FZS 不建 Topic;统一 Topic 模式只需一个,按表分 Topic 模式需要为每张表建。
  2. 分区固定为 0:Topic 分区数建议 1,消费并行度请从下游应用层面解决。
  3. 无事务边界:消息中不含事务开始/结束与 XID,需事务一致性请自行分组或走库表链路。
  4. 消息可能被丢弃:消息体超过 130,000,000 字节、或重试 3 次仍失败的消息会被跳过,请监控 FZS 日志中的 producev fail / too large 关键字。
  5. 空串 = NULL:若下游对空串敏感,需要额外的数据清洗策略。
  6. 二进制为 hex:BLOB / RAW 列需要 hex 解码后再使用。
  7. KZS_JSON 的字段不一致:多行插入用 SCNTIME,单行插入与删除的 message 内没有 OP,解析时请按 OP 分支处理。
  8. DEBEZIUM_JSON 非标准 Debezium:没有 schema 描述,snapshot 是字符串,快照 op 为 c,不要直接复用强校验的 Debezium 解析器。
  9. 时间戳精度到秒:ts_ms 末 3 位为 000,起始阶段可能为 0。