- 1 源库配置
- 1.1 发布选择kafka类型
- 1.2 配置说明
- 1.2.1 断点配置说明
- 1.2.2 配置连接信息
- 1.3 设置 JSON 模板(json_template)
- 1.3.1 配置模板示例
- 1.3.2 参数详细说明
- 1.3.2.1 schema_name(Schema 名称)
- 1.3.2.2 table_name(表名称)
- 1.3.2.3 optime(操作时间)
- 1.3.2.4 optype(操作类型)
- 1.3.2.5 columns(列信息)
- 1.3.2.6 列数据来源(insert_value / delete_value / after_value / before_value / ddl_value)
- 1.3.2.7 data_layer(数据实际位置)
- 1.3.3 其他处理规则
- 1.4 支持key-value的json格式
- 1.5 json_template添加参数data_layer
- 1.6 支持嵌套数组json格式
- 1.7 支持嵌套数组json格式,存在多条记录的情况
- 1.8 支持DSG的json格式
- 1.9 SASL认证配置
- 1.10 支持指定源端topic多分区:partition_count
可以利用kafka工具查看topic中的数据情况
kafka-console-consumer.sh --bootstrap-server ip:port --topic topic名 --from-beginning
例如:
[root@node ~]# docker exec -it zcbuskafka bash
[root@zcbuskafka kafka]# cd bin
[root@zcbuskafka bin]# ./kafka-console-consumer.sh --bootstrap-server zcbuskafka:9092 --topic test001 --from-beginning
{"op_time":"2024-06-04 10:21:20","db_type":"mysql","schema_name":"testdb","table_name":"test","optype":"insert","values":{"id":"257","name":"sdfrsf","optime":"2023-11-05 07:31:13","col1":"2023-11-05","col2":"zbcsdkseebfbds","col3":"850226","col4":"a1AAieFJMAZA","col5":"4b4DLPNyiOFBpcN0imgmh"}}
{"op_time":"2024-06-04 10:21:20","db_type":"mysql","schema_name":"testdb","table_name":"test","optype":"insert","values":{"id":"258","name":"sdfrsf","optime":"2023-11-05 07:31:13","col1":"2023-11-05","col2":"zbcsdkseebfbds","col3":"850226","col4":"a1AAieFJMAZA","col5":"4b4DLPNyiOFBpcN0imgmh"}}
{"op_time":"2024-06-04 10:21:20","db_type":"mysql","schema_name":"testdb","table_name":"test","optype":"insert","values":{"id":"259","name":"sdfrsf","optime":"2023-11-05 07:31:13","col1":"2023-11-05","col2":"zbcsdkseebfbds","col3":"850226","col4":"a1AAieFJMAZA","col5":"4b4DLPNyiOFBpcN0imgmh"}}
{"op_time":"2024-06-04 10:21:20","db_type":"mysql","schema_name":"testdb","table_name":"test","optype":"insert","values":{"id":"260","name":"sdfrsf","optime":"2023-11-05 07:31:13","col1":"2023-11-05","col2":"zbcsdkseebfbds","col3":"850226","col4":"a1AAieFJMAZA","col5":"4b4DLPNyiOFBpcN0imgmh"}}
……1 源库配置
1.1 发布选择kafka类型

1.2 配置说明
1.2.1 断点配置说明

偏移量格式说明:
- 留空:从最早开始消费
- -1:从最新开始消费
- 分区号:偏移量(如 0:19):从指定分区的指定位置开始消费
1.2.2 配置连接信息
对于填写连接信息时 只需填写topic名以及kafka_bootstrap.servers(kafka连接信息),其他参数默认即可。
| 序号 | 配置项 | 配置参数 | 描述信息 |
|---|---|---|---|
| 1 | db_type | KAFKA | 源数据库类型 |
| 2 | data_format | json | kafka中数据的格式,支持bsd/json,默认按照bsd类型解析 |
| 3 | label | - | 标签 |
| 4 | topic | topic名字 | kafka增量数据来源的topic名字 |
| 5 | kafka_auto.commit.interval.ms | 60000 | 自动提交间隔。范围:[0,Integer.MAX],默认值是 5000 (5 s) 消费者偏移的频率以毫秒为单位自动提交给Kafka,如果enable.auto.commit设置为true。 |
| 6 | kafka_auto.offset.reset | earliest | 这个配置项,是告诉Kafka Broker在发现kafka在没有初始offset,或者当前的offset是一个不存在的值(如果一个record被删除,就肯定不存在了)时,该如何处理。它有4种处理方式:? ● earliest: 自动将偏移量重置为最早的偏移量 ● latest:自动将偏移量重置为最新的偏移量 ● none: 如果边更早的offset也没有的话,就抛出异常给consumer,告诉consumer在整个consumer group中都没有发现有这样的offset。 ● anything else: 如果不是上述3种,只抛出异常给consumer。 |
| 7 | kafka_bootstrap.servers | 172.17.58.146:9092 | 用于建立与kafka集群连接的host/port组。数据将会在所有servers上均衡加载,不管哪些server是指定用于bootstrapping。这个列表仅仅影响初始化的hosts(用于发现全部的servers)。这个列表格式: host1:port1,host2:port2,… 因为这些server仅仅是用于初始化的连接,以发现集群所有成员关系(可能会动态的变化),这个列表不需要包含所有的servers(你可能想要不止一个server,尽管这样,可能某个server宕机了)。如果没有server在这个列表出现,则发送数据会一直失败,直到列表可用。 list类型 |
| 8 | kafka_session.timeout.ms | 100000 | 使用Kafka的组管理设施时,用于检测消费者失败的超时。消费者定期发送心跳来向经纪人表明其活跃度。如果代理在该会话超时到期之前没有收到心跳,那么代理将从该组中删除该消费者并启动重新平衡。请注意,该值必须在允许的范围内 Consumer session 过期时间。这个值必须设置在broker configuration中的group.min.session.timeout.ms 与 group.max.session.timeout.ms之间。 |
| 9 | kafka_enable.auto.commit | false | enable.auto.commit 参数用于指定是否自动提交消费者的偏移量。如果设置为 true,消费者会定期自动提交偏移量;如果设置为 false,您需要手动提交偏移量。默认值为 true。 |
| 10 | source_db_type | mysql | 源 类型 |
| 11 | source_db_name | kafka | 实例名【kafka可以忽略】 |
| 12 | source_host | 127.0.0.1 | 数据库主机地址 |
| 13 | source_port | 3306 | 数据库端口号 |
| 14 | source_user | zcbus | 数据库用户名 |
| 15 | source_password | 1qaz!QAZ | 数据库密码 |
| 16 | real_send_cache_buffer | 0 | 增量发布,设置cache中缓存buffer大小,单位M,最小50,最大500,如果设置这个值,real_send_queues忽略 |
| 17 | statistics_into_influxdb | 0 | 是否将增量信息记录到influxdb中 0 不添加 1添加 |
| 18 | send_sys_time | 0 | 发送系统时间的时间间隔,单位为分钟,设置大于0时,每隔指定的时间间隔发送一次系统时间 |
| 19 | send_log_position | 0 | 发送增量分析日志点的时间间隔,单位为分钟,默认为0不发送,设置大于0时,每隔指定的时间间隔发送一次增量分析的日志点、日志时间、系统时间 |
| 20 | max_conf_db_connection | 1 | zcbus配置库的最大连接数,默认是2 |
1.3 设置 JSON 模板(json_template)
1.3.1 配置模板示例
以下为一个完整的 JSON 配置示例:
{
"schema_name": "sc01",
"table_name": [
{ "name": "table_tag1" },
"_",
{ "name": "table_tag2", "concat": "_" },
"_abc"
],
"optime": {
"name": "loadertime",
"convert": "timestamp_to_date"
},
"optype": {
"name": "optype_tag",
"opmap": {
"I": "insert",
"U": "update",
"D": "delete",
"DDL": "ddl"
}
},
"columns": [
{ "name": "col1", "new_name": "col1_new", "convert": "timestamp_ms_to_date" },
{ "name": "col2", "new_name": "col2_new" }
]
}1.3.2 参数详细说明
1.3.2.1 schema_name(Schema 名称)
- 配置
schema_name字段指定目标 Schema。 - 示例:若 Schema 名为
s01,则配置为:"schema_name": "sc01"
1.3.2.2 table_name(表名称)
支持两种方式:
(1)固定表名
直接使用字符串指定:
"table_name": "tb01"(2)动态拼接表名
通过数组形式,结合用户 JSON 中的多个标签值拼接:
"table_name": [
{ "name": "table_tag1" },
"_",
{ "name": "table_tag2", "concat": "_" },
"_abc"
]- 若
table_tag1="ss",table_tag2=["tt1", "tt2"],且concat为"_",则最终表名为:ss_tt1_tt2_abc。 - 若标签位于嵌套对象中,可用
::分隔路径,如obj::obj2::tag。
💡 与发布配置表(bus_push_sync_tb)的对应关系:
schema_name和table_name提取的值 → 关联bus_push_sync_tb配置表 → 确定目标 Topic 和推送规则
1.3.2.3 optime(操作时间)
支持以下三种方式:
(1)默认(不配置)
使用 Kafka 消息中的时间戳。
(2)固定时间
直接指定时间字符串:
"optime": "2020-02-05 10:21:32"(3)从 JSON 字段取值
"optime": {
"name": "loadertime"
}- 若字段位于嵌套对象中,支持
::路径,如obj::obj2::loadertime。 - 若值为时间戳(秒或毫秒),可通过
convert转换:"optime": { "name": "loadertime", "convert": "timestamp_to_date" }
支持的转换函数:
| 函数名 | 说明 |
|---|---|
timestamp_to_date |
将秒级时间戳转为时间字符串 |
timestamp_ms_to_date |
将毫秒级时间戳转为时间字符串 |
timestamp_str_adjust |
校正时间字符串为 BSD 格式(替换 T 为空格,去除毫秒) |
1.3.2.4 optype(操作类型)
支持以下方式:
(1)固定操作类型
"optype": "insert"(2)从 JSON 字段取值
"optype": {
"name": "optype_tag"
}- 字段值需为程序支持的默认类型:
insert、update、delete、ddl。 - 支持嵌套路径,如
obj::obj2::optype_tag。
(3)自定义映射
若字段值非默认类型,可通过 opmap 映射:
"optype": {
"name": "optype_tag",
"opmap": {
"I": "insert",
"U": "update",
"D": "delete",
"DDL": "ddl"
}
}1.3.2.5 columns(列信息)
- 若无需修改列信息,可不配置
columns。 - 若需修改列名或对列值进行转换,配置如下:
"columns": [ { "name": "col1", "new_name": "col1_new", "convert": "timestamp_ms_to_date" }, { "name": "col2", "new_name": "col2_new" } ]
1.3.2.6 列数据来源(insert_value / delete_value / after_value / before_value / ddl_value)
用于指定不同操作类型的数据来源字段,支持嵌套路径(:: 分隔)。
| 配置项 | 说明 | 默认值 |
|---|---|---|
insert_value |
insert 操作的列数据来源 | "values" |
delete_value |
delete 操作的列数据来源 | "values" |
after_value |
update 操作中 after 值的来源 | "after" |
before_value |
update 操作中 before 值的来源 | "before" |
ddl_value |
ddl 操作语句的来源 | "values" |
- 若数据仅为一层 JSON,可设置为
"/"。 - 若为嵌套结构,使用
::,如obj::obj2::tag。
示例:
"insert_value": "values_tag",
"after_value": "obj::after_data"1.3.2.7 data_layer(数据实际位置)
若原始 JSON 外层还有包装,实际有效数据位于第二层,可通过 data_layer 指定解析层数:
"data_layer": 21.3.3 其他处理规则
- 列名中的点号(
.) :自动替换为下划线(_)。 - 数组类型列值:将数组内各成员值用逗号拼接为一个字符串。
1.4 支持key-value的json格式
zcbus-8.3-16-20250415.tar.gz版本处理
json_template参数配置示例:
{
"insert_value": "/",
"schema_name": "testdb",
"table_name": "test",
"optype": "insert"
}兼容数据示例:
{"cycleType":"5m","dataKey":"aep_aiot_api_current_count_2000039548","reportDate":"2024-11-28 00:00:00","provice":"91000000000000001","dataValue":"0.8","type":"0","dataCode":"RZBZ00002261","cycle":"2024-11-28 00:00:00"}1.5 json_template添加参数data_layer
zcbus-8.3-16-20250522.tar.gz版本处理
参数说明:指定第n层为真正需解析的数据
json_template参数配置示例:
{
"insert_value": "/",
"schema_name": "testdb",
"table_name": "test",
"optype": "insert",
"data_layer": 2
}兼容数据示例:
过滤第一层”NULL100989408”标签数据
{ "_NULL_100989408" : {"town":"","city":"","proviceName":"安徽","FromSys":"集团-综合调度系统","dataValue":"0","type":"0","dataCode":"MDXT00002062","cycle":"2025-04","cycleType":"M","dataKey":"completion_count_of_major_hazard_rectification","reportDate":"2025-05-14 09:28:09","district":"","provice":"911010000000000000000000"} }1.6 支持嵌套数组json格式
zcbus-8.3-16-20260129.tar.gz版本处理
(1)data_format设置:json_string
(2)json_template设置:
{
"insert_value": "/",
"schema_name": "testdb",
"table_name": "test",
"optype": "insert"
}兼容数据示例:
[{"metricEngName":"idcRecoveryAmounts","resourceSys":"手工填报","result":0,"creditTerm":"2025-11","province":611000000000001824854982,"dimeRelaList":[{"dimeName":"区域","dimeCode":"AreaDime","dimeValue":"611000000000001824854982","dimeValueCn":"陕西"}]}]1.7 支持嵌套数组json格式,存在多条记录的情况
zcbus-8.3-16-20260421.tar.gz版本处理
(1)data_format设置:json
(2)json_template设置:
{
"insert_value": "/",
"schema_name": "testdb",
"table_name": "test",
"optype": "insert"
}兼容数据示例:
[{"metricEngName":"idcCheckRealNums","resourceSys":"集团IDC+运营子系统","result":60,"creditTerm":"2026-02","province":911010000000000000000000,"dimeRelaList":[{"dimeName":"区域","dimeCode":"AreaDime","dimeValue":"911010000000000000000000","dimeValueCn":"全国"}]},{"metricEngName":"idcCheckRealNums","resourceSys":"集团IDC+运营子系统","result":9,"creditTerm":"2026-02","province":311115000000000000000001,"dimeRelaList":[{"dimeName":"区域","dimeCode":"AreaDime","dimeValue":"311115000000000000000001","dimeValueCn":"上海"}]}]1.8 支持DSG的json格式
zcbus-8.3-16-20260909.tar.gz版本处理
(1)data_format设置:json
(2)json_template设置:
{
"data_layer": 1,
"after_value": "message::data",
"before_value":"message::beforeData",
"insert_value": "message::data",
"delete_value": "message::data",
"ddl_value": "values",
"schema_name": "test",
"table_name": "test",
"optime": {
"name": "message::headers::timestamp",
"convert":"timestamp_str_adjust"
},
"optype": {
"name": "message::headers::operation",
"opmap": {
"INSERT": "insert",
"UPDATE": "update",
"DELETE": "delete",
"DDL": "ddl"
}
}
}兼容数据示例:
insert:
{"magic": "atMSG","type": "DT","headers": null,"messageSchemaId": null,"messageSchema": null,"message": {"data": {"fphm": "26442000010163137351","kprq": "2026-09-04 09:33:23.000000","kjhzfpdydlzfphm": "","fppz_dm": "02","tdyslx_dm": "","sfzlzfp": "Y","qy_dm": "","fpkjfs_dm": "1","xsfnrsbh": "91440101MA9UTX4G5L"},"beforeData": null,"headers": {"operation": "INSERT","changeSequence": "1448199119","timestamp": "2026-09-04T09:33:46.450000","streamPosition": "00000ddc.023a65c3.0000001.0000.02.0001:24444.271273.16","transactionId": "00000000000000000000000000006B0002","changeMask": "0FFFFFFFFFFFFF","columnMask": "0FFFFFFFFFFFFF","transactionEventCounter": 1,"transactionLastEvent": false}}}
update:
{"magic": "atMSG","type": "DT","headers": null,"messageSchemaId": null,"messageSchema": null,"message": { "data": {"fphm": "2692200001036028311","kprq": "2026-09-04 09:33:29.000000","kjhzfpdydlzfphm": "","fppz_dm": "01","tdyslx_dm": "","sfzlzfp": "Y","qy_dm": "","fpkjfs_dm": "1","xsfnrsbh": "91370285MA3LY9E785" }, "SJCSBS": "I", "beforeData": {"fphm": "2692200001036028311","kprq": "2026-09-04 09:33:29.000000","kjhzfpdydlzfphm": "","fppz_dm": "01","tdyslx_dm": "","sfzlzfp": "Y","qy_dm": "","fpkjfs_dm": "1","xsfnrsbh": "91370285MA3LY9E785" }, "headers": {"operation": "UPDATE","changeSequence": "1448199153","timestamp": "2026-09-04T09:33:46.485000","streamPosition": "00000ddc.023a65c3.0000001.0000.02.0001:24444.271273.16","transactionId": "00000000000000000000000000006B0002","changeMask": "0FFFFFFFFFFFFF","columnMask": "0FFFFFFFFFFFFF","transactionEventCounter": 1,"transactionLastEvent": false }}}
delete:
{"magic": "atMSG","type": "DT","headers": null,"messageSchemaId": null,"messageSchema": null,"message": {"data": {"fphm": "26442000010163137351","kprq": "2026-09-04 09:33:23.000000","kjhzfpdydlzfphm": "","fppz_dm": "02","tdyslx_dm": "","sfzlzfp": "Y","qy_dm": "","fpkjfs_dm": "1","xsfnrsbh": "91440101MA9UTX4G5L"},"beforeData": null,"headers": {"operation": "DELETE","changeSequence": "1646734131","timestamp": "2026-09-04T09:42:30.762000","streamPosition": "00000dc023a65c3.00000001.0000.02.0001:24444.271273.16","transactionId": "0000000000000000000006B0002","changeMask": "0xFFFFFFFFFFFF","columnMask": "0xFFFFFFFFFFFF","transactionEventCounter": 1,"transactionLastEvent": false}}}1.9 SASL认证配置

1.10 支持指定源端topic多分区:partition_count

注意 :
(1) kafka发布 暂时可以识别增量数据内容 ;
(2) 发布时,选择批量导入方式 填写库名.表名;
(3) 订阅到目标数据库时 需手动建表,如有主键则使用主键 没有则需指定apply_by_key参数;
