可以利用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": 2

1.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参数;

文档更新时间: 2026-09-11 01:29   作者:程少波