帮你快速理解、总结文档立即下载

Kafka 数据源

最近更新时间:2026-09-11 20:42:31
我的收藏

概述

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')

相关文档