下面是本次完整配置流程,重点包括 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) 输出。

六、部署顺序

  1. 准备客户端 krb5.conf
  2. 部署客户端 Kafka keytab;
  3. 验证 kinitkvno
  4. 配置 Kafka SASL 参数;
  5. 发布端开启 send_sys_time=1
  6. 更新 zcbus_other
  7. 配置对象的 chk_topic
  8. 确认 nodeId SQL 查询结果;
  9. 停止旧 API/client;
  10. 替换最新 JAR;
  11. 启动 API/client;
  12. 发送测试业务数据;
  13. 等待增量时间标记或全量结束标记;
  14. 检查 chk_topic 是否收到稽核消息;
  15. 核对 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名称:分区号
文档更新时间: 2026-08-18 17:48   作者:周风磊