下面是本次完整配置流程,重点包括 Kafka Kerberos 认证、北京移动稽核 Topic、发布端 send_sys_time 以及完整处理链路。
一、Kafka Kerberos 认证配置
1. 服务端 Principal
Kafka Broker 使用:
kafka/zcbus60@EXAMPLE.COM
对应服务端 keytab:
/etc/security/keytabs/kafka.service.keytab
检查:
klist -kte /etc/security/keytabs/kafka.service.keytab
应能看到:
kafka/zcbus60@EXAMPLE.COM
2. 客户端 Principal
客户端使用:
kafka-client@EXAMPLE.COM
客户端 keytab:
/usr/local/zcbus/config/kafka.client.keytab
检查:
klist -kte /usr/local/zcbus/config/kafka.client.keytab
3. 客户端 krb5.conf
例如:
[logging]
default = FILE:/var/log/krb5libs.log
kdc = FILE:/var/log/krb5kdc.log
admin_server = FILE:/var/log/kadmind.log
[libdefaults]
default_realm = EXAMPLE.COM
dns_lookup_realm = false
dns_lookup_kdc = false
rdns = false
ticket_lifetime = 24h
renew_lifetime = 7d
forwardable = true
default_tkt_enctypes = aes128-cts
default_tgs_enctypes = aes128-cts
[realms]
EXAMPLE.COM = {
kdc = zcbus60
admin_server = zcbus60
}
[domain_realm]
.zcbus60 = EXAMPLE.COM
zcbus60 = EXAMPLE.COM
远程客户端必须能够解析:
zcbus60
建议配置 /etc/hosts 或 DNS:
192.168.2.60 zcbus60
4. Kafka 客户端配置
bootstrap.servers=zcbus60:9092
security.protocol=SASL_PLAINTEXT
sasl.mechanism=GSSAPI
sasl.kerberos.service.name=kafka
sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required \
useKeyTab=true \
doNotPrompt=true \
keyTab="/usr/local/zcbus/config/kafka.client.keytab" \
storeKey=true \
useTicketCache=false \
principal="kafka-client@EXAMPLE.COM";
说明:
sasl.kerberos.service.name=kafka
不是客户端用户名,而是 Kafka Broker 服务 Principal 中 / 前面的服务名:
kafka/zcbus60@EXAMPLE.COM
因此这里必须是:
kafka
5. Kerberos 验证
export KRB5_CONFIG=/usr/local/zcbus/config/krb5.conf
export KRB5CCNAME=FILE:/tmp/krb5cc_$(id -u)
kinit -kt /usr/local/zcbus/config/kafka.client.keytab \
kafka-client@EXAMPLE.COM
klist
kvno kafka/zcbus60@EXAMPLE.COM
验证成功后再启动 API/client。
如果出现:
Server not found in Kerberos database
重点检查:
bootstrap.servers是否使用zcbus60,不要直接使用 IP;- Broker Principal 是否为 zcbus60@EXAMPLE.COM"">`kafka/zcbus60@EXAMPLE.COM`;
- 主机名是否能解析;
sasl.kerberos.service.name是否为kafka;- keytab 中 Principal 是否正确。
二、稽核 Topic 配置
1. 发布端必须开启 send_sys_time
增量稽核消息依赖上游发送周期时间标记。
发布端需要开启:
send_sys_time=1
该参数用于发送系统时间或周期结束标记。没有该标记时,客户端无法判断增量统计周期结束,稽核消息可能一直不生成。
发布端还必须正确发送全量标记:
BSDATA_TAG_FULL_SYNC_START
BSDATA_TAG_FULL_SYNC_END
全量稽核消息在收到 FULL_SYNC_END 后生成。
2. zcbus_other 配置
北京移动推荐配置:
{
"jsontype": "json_module",
"keymap": {
"optype": "operType",
"op_time": "opTs",
"loaderTime": "loaderTime",
"batchCode": "batchCode",
"table_name": "tableName",
"schema_name": "tableOwner",
"before": "columnInfoBefore",
"after": "columnInfoAfter",
"values": "columnInfoAfter",
"delete": "columnInfoBefore"
},
"opermap": {
"fullload": "I",
"insert": "I",
"update": "U",
"delete": "D",
"ddl": "DDL"
},
"module_check": "bjyd_check_v1",
"real_check": 1,
"full_check": 1,
"sendtype": 1,
"ifIncludeCheckZeroData": 0,
"columnCompleteMode": "all",
"columnOrder": "table_def_order",
"nullValueOutput": "json_null",
"beforeOnInsert": "skip",
"beforeOnUpdate": "empty_object",
"afterOnDelete": "skip",
"opFormat": "opermap",
"batchCodeFormat": "yyyyMMddHHmm",
"topLevelFields": [
"loaderTime",
"columnInfoAfter",
"test_Own",
"opTs",
"columnNum",
"tableOwner",
"batchCode",
"columnInfoBefore",
"rid",
"operType",
"tableName"
],
"appendFields": {
"rid": "",
"test_Own": "${tableOwner}",
"tableOwner": "${tableOwner}",
"columnNum": "${columnNum}"
},
"ColNameCase": 2,
"ifTableNameCase": 2,
"ifAllDataToString": 1,
"ifLoaderTimeAsTimestamp": 0,
"ifOpTimeAsTimestamp": 0,
"LoaderTimeFormat": "yyyy-MM-dd HH:mm:ss",
"ifIncludeloaderTime": 1,
"ifIncludeOptype": 1,
"ifIncludebatchCode": 1,
"ifIncludeTableName": 1,
"ifIncludeSchemaName": 1,
"ifIncludeoptime": 1,
"ifInclueUpdateBefore": 1,
"ifIncludeCheckZeroData": 0,
"skipddl": 1
}
重点参数:
module_check=bjyd_check_v1
生成北京移动扁平稽核 JSON。
real_check=1
开启增量稽核。
full_check=1
开启全量稽核。
sendtype=1
实际发送 Kafka。
ifIncludeCheckZeroData=0
当前代码逻辑下,设置为 0 才允许零数据周期发送稽核消息。设置为 1 时,零数据周期会被跳过。
3. 配置 chk_topic
chk_topic 不写在 zcbus_other JSON 中,而是在对象属性表中配置:
insert into bus_client_product_object_properties
(customerid, objectid, variable_name, value, description)
values (
你的customerid,
你的objectid,
'chk_topic',
'chk_bjyd_test_api_order:0',
'check topic for kafka'
);
批量配置示例:
insert into bus_client_product_object_properties
(customerid, objectid, variable_name, value, description)
select customerid,
id,
'chk_topic',
concat('chk_', targettableschema, ':0'),
'check topic for kafka'
from bus_client_product_object
where targettableschema like 'bjyd_test_api_%'
and id not in (
select objectid
from bus_client_product_object_properties
where variable_name = 'chk_topic'
);
格式:
Topic名称:分区号
例如:
chk_bjyd_test_api_order:0
如果不带分区号,程序默认使用 0 分区,但建议明确配置。
三、关键配置检查 SQL
检查 zcbus_other:
select dbid, variable_name, value
from bus_client_db_parameter
where variable_name = 'zcbus_other';
检查稽核 Topic:
select customerid, objectid, variable_name, value
from bus_client_product_object_properties
where variable_name in (
'chk_topic',
'chk_real_topic',
'chk_full_topic'
);
检查上游节点 ID:
select bcpo.id as objectid,
bcpo.tableid,
bit.nodeid
from bus_client_product_object bcpo
join bus_in_tables bit
on bit.id = bcpo.tableid
where bcpo.id = 你的objectid;
其中:
bit.nodeid
就是北京移动稽核消息中的:
"nodeId": 10007
四、完整处理逻辑
增量链路
上游产生 INSERT / UPDATE / DELETE
↓
发布端推送业务数据
↓
发布端 send_sys_time=1
↓
周期结束时发送系统时间标记
↓
客户端累计 insertCount/updateCount/deleteCount
↓
收到时间标记
↓
记录 beginTime/endTime
↓
读取上游 nodeId
↓
生成 bjyd_check_v1 扁平 JSON
↓
发送到 chk_topic
↓
清空当前周期统计
↓
开始下一个周期
全量链路
收到 FULL_SYNC_START
↓
初始化全量统计
↓
持续接收全量数据
↓
累计全量记录数
↓
收到 FULL_SYNC_END
↓
记录全量结束时间
↓
生成全量稽核 JSON
↓
发送到 chk_topic
↓
清空全量统计
北京移动稽核消息示例
{
"nodeName": "public",
"updateCount": 151,
"dataRange": "2026-08-19 09:40:00-2026-08-19 09:44:57",
"tableOwner,test_Own": "SO2",
"batchCode": "202608190940",
"ddl": 0,
"tableName": "INS_OFFER_109",
"dataCount": 622,
"loaderTime": "2026-08-19 09:45:19",
"deleteCount": 4,
"tableOwner": "public",
"insertCount": 467,
"beginTime": "2026-08-19 09:40:00",
"endTime": "2026-08-19 09:44:57",
"dsgStatus": 0,
"nodeId": 10007,
"errorCount": 0
}
五、日志核对重点
Kerberos 正常时应能看到认证成功和 Kafka 连接日志。
稽核生成时重点查看:
json_mode : bjyd_check_v1
生成但未发送:
Skip Send check data
实际发送:
Send check data ... to topic ...
节点 ID 查询调试日志:
Get upstream nodeId SQL ...
Get upstream nodeId result ...
如果看不到 SQL 调试日志,需要将日志级别调整为允许 print(2) 输出。
六、部署顺序
- 准备客户端
krb5.conf; - 部署客户端 Kafka keytab;
- 验证
kinit和kvno; - 配置 Kafka SASL 参数;
- 发布端开启
send_sys_time=1; - 更新
zcbus_other; - 配置对象的
chk_topic; - 确认
nodeIdSQL 查询结果; - 停止旧 API/client;
- 替换最新 JAR;
- 启动 API/client;
- 发送测试业务数据;
- 等待增量时间标记或全量结束标记;
- 检查
chk_topic是否收到稽核消息; - 核对
nodeId、批次号和统计数量。
当前最新 JAR:
target\ClientApiMain.jar
SHA-256:
981B11F29A98871F632D1473A2C723EF955222E7AC9C19C1C693DA285C74E8D5
重点需要修改的配置可以归纳为:
Kerberos:
krb5.conf
kafka.client.keytab
sasl.jaas.config
security.protocol
sasl.mechanism
sasl.kerberos.service.name
bootstrap.servers
发布端:
send_sys_time=1
客户端 zcbus_other:
module_check=bjyd_check_v1
real_check=1
full_check=1
sendtype=1
ifIncludeCheckZeroData=0
对象属性:
chk_topic=Topic名称:分区号