概述
Kafka 是分布式流处理平台。DataBuddy 数据接入支持将 Kafka(含腾讯云 CKafka)作为来源端读取消息,覆盖离线读取与实时整库场景。本文按「建数据源 → 配置同步任务 → 处理常见问题」的顺序介绍其使用方式。
支持的版本
类型 | 支持版本 |
自建 Kafka | 2.x / 3.x |
腾讯云 CKafka | 2.4.1 / 2.8.1 / 3.2.3 |
读取能力
能力 | 离线读 | 实时整库读 |
支持 | ✓ | ✓ |
使用限制
序列化格式支持
场景 | 读取格式 |
整库同步 | canal-JSON / debezium / ogg-json / csv /json /raw |
整库同步的 DDL 变更仅支持新增列、删除列、新增表 。
创建数据源
操作步骤
1. 登录 DataBuddy 控制台。
2. 在顶部切换到目标地域和 Workspace。
3. 进入 数据接入 > 任务管理 > 数据源管理 (或 平台管理 > 数据源管理 )。
4. 点击 添加数据源 ,选择 Kafka 。
5. 填写连接配置与认证信息,详见 参数说明。
6. 点击 测试连接 (参见 数据源连通性测试)。
7. 测试通过后点击 保存 或 保存 & 创建任务 。
参数说明
参数 | 说明 | 是否必填 |
数据源名称 | 数据源在工作空间内的唯一标识 | 是 |
Kafka 类型 | 开源自建 Kafka / 腾讯云 CKafka | 是 |
Bootstrap Servers | host1:port1,host2:port2(如 10.0.0.1:9092,10.0.0.2:9092) | 是 |
认证方式 | 见 认证方式 | 否 |
用户名 / 密码 | 选择 SASL 认证时必填(密码支持 SSM 凭证托管) | 视认证 |
认证方式
方式 | 说明 | 必填字段 |
无认证(PLAINTEXT) | 不开启认证 | — |
SASL_PLAINTEXT | 明文 SASL 认证 | 用户名、密码 |
SASL_GSSAPI(Kerberos) | Kerberos 认证 | keytab 文件、conf 文件、principal、Kafka 服务名称 |
SASL_SSL | SASL + SSL 双重认证 | Truststore 认证文件、Keystore 密码、用户名、密码 |
SSL | TLS 双向认证 | Truststore 认证文件、Truststore 密码、Keystore 认证文件、Keystore 密码、私钥口令 |
网络与服务端预备
集成资源组到所有 broker 必须双向连通 。
服务端
server.properties 中的 listeners / advertised.listeners 必须填对外可达 地址(很多连接失败的根因在此)。在数据接入任务中使用
实时同步:来源端
参数 | 说明 |
数据源 | 已创建的 Kafka 数据源,支持连通性测试 |
Topic | 要消费的 Topic,支持下拉多选 |
序列化格式 | 见 序列化格式支持 |
读取位置 | 最早 (earliest)/ 最新 (latest)/ 指定时间 |
订阅类型 | 可勾选 INSERT / UPDATE / DELETE,至少勾选一种 |
嵌套 JSON
默认不支持嵌套 JSON,开启需在高级参数中加:
source.nested.json.enabled=true
提示
该参数同时作用于 Kafka key 和 value 。如果 key 不是 JSON(普通字符串、数字主键),需要再追加
source.key.nested.json.enabled=false,否则 key 解析会报错。引用嵌套字段时,使用
父字段->子字段 形式声明,多层嵌套用多个 -> 连接。在函数中引用必须用反引号包裹:TO_TIMESTAMP(`book->published`, 'yyyyMMdd')
高级参数
除下表所列参数外,所有 Kafka 原生客户端参数加
properties. 前缀均可透传,例如 properties.fetch.max.wait.ms=500。参数 | 说明 | 默认 |
source.nested.json.enabled | 启用嵌套 JSON(同时作用于 key 和 value) | false |
source.key.nested.json.enabled | 是否对 key 启用嵌套 JSON | 跟随 source.nested.json.enabled |
source.ignore-parse-errors | 解析失败时跳过错误记录 | false |
scan.topic-partition-discovery.interval | 周期性发现新 partition / 新 Topic 的间隔;填 0 关闭 | 5min |
timestamp.format | 时间字段解析格式 | ISO8601 |
数据格式样例
canal-json
{"data": [{ "id": "1", "name": "张三", "age": "25" }],"database": "test_db","table": "user_info","type": "INSERT","es": 1589373560000,"ts": 1589373560798,"isDdl": false,"mysqlType": { "id": "int", "name": "varchar(50)", "age": "int" },"sqlType": { "id": 4, "name": 12, "age": 4 },"pkNames": ["id"],"old": null,"sql": ""}
debezium
{"before": null,"after": { "id": 1, "name": "张三", "age": 25 },"source": {"version": "1.5.0.Final","connector": "mysql","name": "mysql_server","ts_ms": 1589373560000,"db": "test_db","table": "user_info"},"op": "c","ts_ms": 1589373560798}
常见问题
Q:连接失败、测试连接超时?
A:按以下顺序排查:
1. 集成资源组到 broker 网络是否双向通(VPC、白名单、安全组)。
2.
advertised.listeners 返回的地址在资源组所在网络是否可达。3. SASL 认证下用户名密码、SASL 机制是否与服务端一致。
Q:消费不到数据?
可能原因 | 怎么确认 |
读取位置选了 latest 但 Topic 没新消息 | 改成 earliest 或指定时间点重启 |
Topic 实际无数据 | 用 Kafka 控制台直接 consume 一次 |
Consumer group 在外部已有持久化 offset | 换一个 group ID 重启 |
「过滤操作」/「库表范围」把数据全过滤掉 | 检查同步任务的对应配置 |
Q:序列化格式解析失败、字段全为 null?
A:按以下顺序排查:
1. 用消费工具拉一条原始消息,对比配置的序列化格式。
2. 消息为嵌套 JSON 但未开启
source.nested.json.enabled=true。3. key 不是 JSON 但启用了嵌套 JSON,需补
source.key.nested.json.enabled=false。Q:整库报 mysqlType / sqlType / pkNames 缺失?
A:
canal-json / debezium 整库需消息携带这些元字段,确认上游 Canal / Debezium 输出包含。用工具抓一条消息核对即可。Q:嵌套字段函数处理报 SQL 语法错误?
A:
- 与 > 不是合法 SQL 标识符。在函数中引用嵌套字段必须用反引号包裹,并确认字段已通过「添加字段」声明:TO_TIMESTAMP(`book->published`, 'yyyyMMdd')