本文为您介绍如何使用钉钉知识库 CDC 连接器。
背景信息
钉钉知识库 CDC 连接器通过周期性轮询钉钉开放平台 OpenAPI,按文档修改时间持续捕获钉钉知识库(Wiki Workspace)中在线文档的新增和更新。连接器输出文档的标题、链接、所属知识库、父节点、更新时间等元数据和 Markdown 格式正文,也可以选择获取文档的原始块结构,适用于 RAG 知识库构建、文档检索索引和湖仓内容同步等场景。
钉钉知识库没有可供直接消费的变更事件流,钉钉知识库 CDC 通过 OpenAPI 轮询实现全量和增量读取:
每个知识库是一个独立的读取分片。未配置
workspace-id时,连接器按scan.discovery.interval周期调用知识库列表接口发现操作人可见的全部知识库;配置后仅同步指定知识库。每次拉取时,连接器按
scan.discovery.max-depth和scan.discovery.max-nodes遍历知识库节点树,筛选出钉钉在线文档节点(type=FILE且category=ALIDOC),并按文档修改时间生成本次拉取的游标位置。游标之后的文档按修改时间排序,单个分片单次最多输出
scan.prefix-size篇。对每篇变化的文档,连接器通过一致性校验机制读取其块结构并渲染成 Markdown 正文。无增量变化时,连接器按
scan.poll-interval间隔重复上述检查。
连接器按主键 doc_id 输出 upsert 语义的 changelog,包含 INSERT 和 UPDATE_AFTER,不输出 UPDATE_BEFORE:
场景 | 输出 |
全量阶段首次读到的文档 | INSERT |
增量阶段内容或元数据变化的文档 | UPDATE_AFTER |
当前实现不感知文档删除:文档被删除后不再出现在节点列表中,连接器不会输出 DELETE 事件,下游中该行会保留为最后一次同步的状态。如需清理残留数据,可提交批模式作业,以下游表与源表按主键做反连接(anti-join),删除主键已不存在于源表中的行。
钉钉 CDC 连接器支持的信息如下。
类别 | 详情 |
支持类型 | 源表、数据摄入数据源 |
运行模式 | 流模式;数据摄入 YAML 另支持批模式(仅 |
数据格式 | 固定 Schema,不支持裁剪字段 |
API 种类 | SQL 和数据摄入 YAML |
同步对象 | 知识库中操作人有权读取的钉钉在线文档节点(节点类型 |
启动模式 |
|
Watermark | 不支持 |
维表和结果表 | 不支持 |
前提条件
网络要求
Flink 作业需要访问钉钉开放平台 OpenAPI(默认域名为 api.dingtalk.com)。作业部署在 VPC 内时,请确保具备公网访问能力(参见阿里云实时计算 Flink 版文档网络连接)。
应用创建与接口授权
无论使用哪种启动模式,都需要先在钉钉开放平台完成应用创建和接口授权。
登录钉钉开发者后台,创建企业内部应用,在应用的"凭证与基础信息"页面获取应用的 Client ID(AppKey)和 Client Secret(AppSecret)。连接器使用 AppKey 和 AppSecret 调用
POST /v1.0/oauth2/accessToken获取企业内部应用 access_token 并自动续期。为应用开通本连接器所需的接口权限。在开发者后台进入目标应用的"权限管理"页面,搜索下表中的权限点并申请开通(部分权限点需要组织管理员审批)。连接器会调用以下钉钉 OpenAPI:
API | 用途 | 权限点 |
| 获取企业内部应用 access_token | 基础能力,无需额外权限点 |
| 列出操作人可见的知识库 | 知识库读权限 |
| 遍历知识库节点树,列出子节点 | 知识库节点读权限 |
| 查询节点元数据 | 知识库节点读权限 |
| 查询文档块元素 | 企业存储文件读权限 |
权限点名称以钉钉开放平台官方文档为准,汇总如下:
能力 | 对应 API | 所需权限 | 是否必需 |
获取访问凭证 |
| 无需额外权限点 | 是 |
读取知识库列表 |
| 知识库读权限 | 是 |
读取节点列表和节点元数据 |
| 知识库节点读权限 | 是 |
读取文档块内容 |
| 企业存储文件读权限 | 是 |
操作人与知识库数据授权
operator-id 填写一名企业成员的 Union ID(可通过钉钉开放平台"查询用户详情"接口获取),连接器所有接口都以该操作人身份调用,可见范围由操作人自身的权限决定:
操作人必须对目标知识库或文档至少拥有"仅可查看"权限("可查看/下载"及以上亦可);
若作业未配置
workspace-id,连接器只能发现操作人可见的知识库。
钉钉知识库支持知识库级、文件夹级和单文档级三层权限,文档默认继承上级目录权限。权限类型及本连接器的最低要求如下:
权限类型 | 能否被本连接器读取 | 说明 |
仅可查看 | 能 | 满足最小读取要求 |
可查看/下载 | 能 | 推荐,避免个别接口对下载类能力的额外要求 |
可编辑 | 能 | 超出读取所需,不建议仅为同步授予 |
可管理 | 能 | 超出读取所需,不建议仅为同步授予 |
无权限 | 不能 | 该知识库或文档不会出现在同步结果中 |
授权方式:由知识库管理员在知识库的"设置 > 成员及权限"中添加操作人(支持按人、部门、群授权),或在单个文档的分享设置中为操作人开放权限。新建文档默认继承知识库权限,按整库授权后新增文档无需单独配置。连接器为只读链路,不需要任何写权限。
授权问题排查
连接器对授权类错误不会重试(会直接失败或跳过对应文档),请按下表排查:
错误码 / 现象 | 含义 | 处理方式 |
| 跨组织访问被拒绝 | 确认应用、操作人和目标知识库属于同一组织;跨组织知识库无法同步 |
| 组织认证等级不满足接口要求 | 在钉钉管理后台完成企业认证或提升认证等级后重试 |
权限点不足(permission denied 类报错) | 应用未开通所需权限点 | 回到"权限管理"页面核对知识库读权限、知识库节点读权限、企业存储文件读权限是否均已开通 |
知识库 / 文档不出现在同步结果中 | 操作人无权读取 | 在知识库"成员及权限"中为操作人添加"仅可查看"及以上权限 |
| 节点已删除或 ID 失效 | 连接器自动跳过,无需处理 |
| 文档标识非法 | 属于永久失败,不会重试,请核对节点状态或联系钉钉支持 |
使用限制
仅支持源表,不支持维表和结果表,不支持 Watermark。
只同步知识库中的钉钉在线文档节点(
type=FILE且category=ALIDOC);文件夹、表格、上传的文件等其他对象会被跳过。当前实现不感知文档删除,不输出 DELETE 事件。
不支持
timestamp启动模式:钉钉 OpenAPI 无法读取文档历史内容。Flink SQL 一张源表只读取一个知识库(
workspace-id也支持逗号分隔多个值,但建议按库建表);数据摄入 YAML 支持逗号分隔多个知识库或不配置自动发现。YAML Schema 与 SQL Schema 不同:YAML 的
blocks_json是受scan.block.include-raw-json控制的物理列,workspace_id为 NOT NULL;请根据接入方式选择对应的 Schema 定义。增量链路基于 OpenAPI 轮询,会产生持续的 API 调用。钉钉开放平台对组织级 API 调用量有配额限制(普通版调用配额较低,专属版配额更高),知识库内文档较多时,请结合
scan.poll-interval、scan.prefix-size和限流参数控制调用成本与新内容的同步延迟。进程级限流参数(
scan.discovery/metadata/content.rate-limit.requests-per-second)由同一 app-key 的所有客户端共享,多个作业共用同一应用时请整体评估配额。content是从文档块结构渲染出的 Markdown 文本,不保证完整保留所有富媒体和结构化元素(无法渲染的块会以占位注释形式标注);需要原始块结构时,Flink SQL 声明blocks_jsonMETADATA 列,数据摄入 YAML 开启scan.block.include-raw-json输出blocks_json物理列。故障恢复可能重复输出 INSERT 和 UPDATE_AFTER,建议下游按主键
doc_id幂等写入。资源发现是累加的:知识库或文档一旦被作业发现,即使之后调小
workspace-id或追加scan.document.exclude-node-ids,也不会从已有作业中移除。需要彻底移除时,请停止作业并无状态重新启动。授权类错误(
forbidden.acrossOrg、orgAuthLevelNotEnough、非法 docKey 等)属于永久失败,连接器不会自动重试,请参见"授权问题排查"处理。
SQL
特色功能
固定 Schema 源表:输出文档标题、链接、正文(Markdown)和元数据,按主键
doc_id输出 upsert 语义的 changelog。启动模式:通过
scan.startup.mode选择读取方式。
模式 | 是否读取全量 | 是否读取增量 | 首次启动行为 |
| 是 | 是 | 先全量读取范围内的文档,再持续轮询增量变化 |
| 是 | 否 | 只全量读取范围内的文档,完成后结束 |
| 否 | 是 | 只处理作业启动后发生的文档变化 |
原始块结构输出:通过
blocks_jsonMETADATA 列输出文档原始块结构的 JSON 文本,需同时开启scan.block.include-raw-json。
语法结构
Flink SQL DDL 必须声明下表中全部 7 个物理字段,字段顺序、名称和大小写需保持一致(字段名大小写敏感),并声明 PRIMARY KEY (doc_id) NOT ENFORCED。
DDL 字段 | 类型 | 说明 |
| STRING NOT NULL | 文档节点标识,源表主键 |
| STRING | 文档所属知识库 ID |
| STRING | 文档在知识库节点树中的父节点 ID |
| STRING | 文档标题 |
| STRING | 文档访问链接 |
| STRING | 文档正文,由文档块结构渲染成的 Markdown 文本 |
| TIMESTAMP_LTZ(3) | 文档最后修改时间 |
此外,DDL 可以通过 METADATA 声明以下只读元数据。元数据列不改变上述物理 Schema。
Metadata Key | 类型 | 说明 |
| STRING | 文档原始块结构的 JSON 文本,可用于保留块级结构或自定义解析。声明该列时必须同时设置 |
CREATE TEMPORARY TABLE dingtalk_wiki_source (
doc_id STRING NOT NULL,
workspace_id STRING,
parent_node_id STRING,
title STRING,
url STRING,
content STRING,
modified_time TIMESTAMP_LTZ(3),
blocks_json STRING METADATA VIRTUAL,
PRIMARY KEY (doc_id) NOT ENFORCED
) WITH (
'connector' = 'dingtalk-cdc',
'object-type' = 'document',
'app-key' = '<yourAppKey>',
'app-secret' = '<yourAppSecret>',
'operator-id' = '<yourOperatorUnionId>',
'workspace-id' = '<yourWorkspaceId>',
'scan.startup.mode' = 'initial',
'scan.block.include-raw-json' = 'true'
);WITH参数
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
connector | 表类型。 | String | 是 | 无 | 固定值为 |
object-type | 同步的对象类型。 | String | 是 | 无 | 当前仅支持 |
app-key | 钉钉企业内部应用的 Client ID(AppKey),用于获取 access_token。 | String | 是 | 无 | 无。 |
app-secret | 钉钉企业内部应用的 Client Secret(AppSecret)。 | String | 是 | 无 | 无。 |
operator-id | 操作人的 Union ID,该用户必须有权读取待同步的文档。 | String | 是 | 无 | 可通过钉钉开放平台"查询用户详情"接口获取。 |
endpoint | 钉钉开放平台 OpenAPI 域名。 | String | 否 |
| 无。 |
workspace-id | 将发现范围限定到指定知识库。 | String | 否 | 无 | 支持逗号分隔多个值。 说明 不配置时自动发现操作人可见的全部知识库。建议显示配置 |
scan.document.exclude-node-ids | 排除的节点 ID 列表。 | String | 否 | 无 | 逗号分隔,这些节点在发现和拉取阶段都会被跳过。 |
scan.discovery.max-depth | 单次知识库轮询扫描允许的最大节点树深度。 | Integer | 否 | 100 | 无。 |
scan.discovery.max-nodes | 单次知识库轮询扫描允许访问的最大节点数。 | Integer | 否 | 100000 | 无。 |
scan.startup.mode | 启动模式。 | String | 否 | initial | 支持 |
scan.discovery.interval | 知识库分片的发现间隔。 | Duration | 否 | 1min | 无。 |
scan.poll-interval | 无增量变化时的轮询间隔。 | Duration | 否 | 10s | 每次轮询遍历知识库节点元数据,仅对按修改时间排序后的变化文档前缀读取正文。 |
scan.prefix-size | 单个知识库单次扫描最多输出正文的文档数。 | Integer | 否 | 100 | 无。 |
scan.document.consistency.max-retries | 获取元数据稳定的文档读取的最大重试次数。 | Integer | 否 | 3 | 设为 0 关闭重试。 |
scan.document.fetch.max-duration | 单篇文档一次一致性内容拉取允许的最长耗时。 | Duration | 否 | 10min | 无。 |
scan.document.max-bytes | 单篇文档允许的最大字节数(含渲染后的 Markdown 正文)。 | Long | 否 | 52428800(50 MiB) | 无。 |
scan.block.fetch-size | 单次 Block API 调用请求的顶层块数量。 | Integer | 否 | 10 | 无。 |
scan.block.max-pages | 单篇文档最多读取的 Block API 分页数。 | Integer | 否 | 1000 | 无。 |
scan.block.include-raw-json | 是否通过 | Boolean | 否 | false | 声明 |
request.timeout | 单次 OpenAPI 请求超时时间。 | Duration | 否 | 30s | 无。 |
request.max-retries | 单次 OpenAPI 请求失败的最大重试次数。 | Integer | 否 | 3 | 无。 |
request.retry-after.max-delay | 可接受的 Retry-After 退避时长上限。 | Duration | 否 | 60s | 超过该值时请求按失败处理。 |
scan.discovery.rate-limit.requests-per-second | 发现类请求(列出知识库和节点)的进程级限流。 | Double | 否 | 5 | 同一 app-key 的所有客户端共享。 |
scan.metadata.rate-limit.requests-per-second | 节点元数据请求的进程级限流。 | Double | 否 | 20 | 同一 app-key 的所有客户端共享。 |
scan.content.rate-limit.requests-per-second | Block API 请求的进程级限流。 | Double | 否 | 10 | 同一 app-key 的所有客户端共享。 |
scan.fetch.rate-limit.requests-per-second | 全局拉取限流。 | Double | 否 | 10 | 在 Source 全部子任务之间分配。 |
代码示例
全量加增量同步知识库文档并写入 print 结果表:
CREATE TEMPORARY TABLE dingtalk_wiki_source (
doc_id STRING NOT NULL,
workspace_id STRING,
parent_node_id STRING,
title STRING,
url STRING,
content STRING,
modified_time TIMESTAMP_LTZ(3),
blocks_json STRING METADATA VIRTUAL,
PRIMARY KEY (doc_id) NOT ENFORCED
) WITH (
'connector' = 'dingtalk-cdc',
'object-type' = 'document',
'app-key' = '<yourAppKey>',
'app-secret' = '<yourAppSecret>',
'operator-id' = '<yourOperatorUnionId>',
'workspace-id' = '<yourWorkspaceId>',
'scan.startup.mode' = 'initial',
'scan.block.include-raw-json' = 'true'
);
CREATE TEMPORARY TABLE paimon_sink (
doc_id STRING,
workspace_id STRING,
parent_node_id STRING,
title STRING,
url STRING,
content STRING,
modified_time TIMESTAMP_LTZ(3)
) WITH (
'connector' = 'paimon',
...
);
INSERT INTO paimon_sink
SELECT doc_id, workspace_id, parent_node_id, title, url, content, modified_time
FROM dingtalk_wiki_source;snapshot 只做一次全量读取,适合内容迁移或一次性导出:
CREATE TEMPORARY TABLE dingtalk_wiki_snapshot (
doc_id STRING NOT NULL,
workspace_id STRING,
parent_node_id STRING,
title STRING,
url STRING,
content STRING,
modified_time TIMESTAMP_LTZ(3),
PRIMARY KEY (doc_id) NOT ENFORCED
) WITH (
'connector' = 'dingtalk-cdc',
'object-type' = 'document',
'app-key' = '<yourAppKey>',
'app-secret' = '<yourAppSecret>',
'operator-id' = '<yourOperatorUnionId>',
'workspace-id' = '<yourWorkspaceId>',
'scan.startup.mode' = 'snapshot'
);如需排除个别文档,在 WITH 参数中追加 'scan.document.exclude-node-ids' = '<nodeId1>,<nodeId2>'。
数据摄入
使用钉钉 CDC Pipeline 连接器,您可以将钉钉知识库中在线文档的全量和增量变化持续同步至下游数据湖或数据仓库。数据摄入连接器的 type 固定为 dingtalk-cdc。
特色功能
自动建表:每个知识库对应一张名为
<知识库ID>.documents的表,连接器在首次发现知识库时自动创建表并同步 Schema。多知识库同步:
workspace-id支持逗号分隔多个值;不配置时自动发现操作人可见的全部知识库,每个知识库自动路由到各自的表。原始块结构输出:开启
scan.block.include-raw-json: true后,输出 Schema 中包含blocks_json物理列。
注意事项
批模式(BATCH)作业仅支持
snapshot启动模式。全量阶段(snapshot 分片)输出 INSERT 事件;增量阶段输出 REPLACE(upsert)事件,下游按主键
doc_id幂等写入即可。运行过程中新发现的知识库固定按
initial方式读取:先读取当前完整内容,再跟随后续变化。作业从 checkpoint/savepoint 恢复时不允许变更
scan.block.include-raw-json,否则作业会报错拒绝启动。仅支持钉钉在线文档节点(
type=FILE且category=ALIDOC),其他对象会被跳过;不感知文档删除。
语法结构
YAML 固定 Schema 如下,主键为 doc_id。与 Flink SQL 相比,YAML Schema 中 blocks_json 作为物理列直接输出(而非 METADATA 列),且 workspace_id 为 NOT NULL。
字段 | 类型 | 说明 |
| STRING NOT NULL | 文档节点标识,主键 |
| STRING NOT NULL | 文档所属知识库 ID |
| STRING | 文档在知识库节点树中的父节点 ID |
| STRING | 文档标题 |
| STRING | 文档访问链接 |
| STRING | 文档正文,由文档块结构渲染成的 Markdown 文本 |
| STRING | 文档原始块结构的 JSON 文本,仅当 |
| TIMESTAMP_LTZ(3) | 文档最后修改时间 |
source:
type: dingtalk-cdc
name: DingTalk Wiki Source
object-type: document
app-key: <yourAppKey>
app-secret: <yourAppSecret>
operator-id: <yourOperatorUnionId>
workspace-id: <yourWorkspaceId>
scan.startup.mode: initial
route:
- source-table: <yourWorkspaceId>.documents
sink-table: dingtalk_wiki.documents
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: <yourDlfCatalogEndpoint>
catalog.properties.warehouse: <yourWarehouse>以上示例把指定知识库中的文档全量读取后持续同步增量变化,写入 Paimon。同步多个知识库时,workspace-id 填写逗号分隔的多个 ID;需要保留原始块结构时追加 scan.block.include-raw-json: true。
配置项
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
type | 连接器类型。 | String | 是 | 无 | 固定值为 |
object-type | 同步的对象类型。 | String | 是 | 无 | 当前仅支持 |
app-key | 钉钉企业内部应用的 Client ID(AppKey)。 | String | 是 | 无 | 无。 |
app-secret | 钉钉企业内部应用的 Client Secret(AppSecret)。 | String | 是 | 无 | 无。 |
operator-id | 操作人的 Union ID,该用户必须有权读取待同步的文档。 | String | 是 | 无 | 无。 |
endpoint | 钉钉开放平台 OpenAPI 域名。 | String | 否 |
| 无。 |
workspace-id | 逗号分隔的知识库 ID 列表。 | String | 否 | 无 | 不配置时自动发现操作人可见的全部知识库。 |
scan.document.exclude-node-ids | 排除的节点 ID 列表。 | String | 否 | 无 | 逗号分隔,这些节点在发现和拉取阶段都会被跳过。 |
scan.discovery.max-depth | 单次知识库轮询扫描允许的最大节点树深度。 | Integer | 否 | 100 | 无。 |
scan.discovery.max-nodes | 单次知识库轮询扫描允许访问的最大节点数。 | Integer | 否 | 100000 | 无。 |
scan.startup.mode | 启动模式。 | String | 否 | initial | 支持 |
scan.discovery.interval | 新知识库的发现间隔。 | Duration | 否 | 1min | 无。 |
scan.poll-interval | 无增量变化时的轮询间隔。 | Duration | 否 | 10s | 无。 |
scan.prefix-size | 单个知识库单次扫描最多输出正文的文档数。 | Integer | 否 | 100 | 无。 |
scan.document.consistency.max-retries | 获取元数据稳定的文档读取的最大重试次数。 | Integer | 否 | 3 | 设为 0 关闭重试。 |
scan.document.fetch.max-duration | 单篇文档一次一致性内容拉取允许的最长耗时。 | Duration | 否 | 10min | 无。 |
scan.document.max-bytes | 单篇文档允许的最大字节数(含渲染后的 Markdown 正文)。 | Long | 否 | 52428800(50 MiB) | 无。 |
scan.block.fetch-size | 单次 Block API 调用请求的顶层块数量。 | Integer | 否 | 10 | 无。 |
scan.block.max-pages | 单篇文档最多读取的 Block API 分页数。 | Integer | 否 | 1000 | 无。 |
scan.block.include-raw-json | 是否在输出 Schema 中包含 | Boolean | 否 | false | 作业恢复时不允许变更该配置。 |
request.timeout | 单次 OpenAPI 请求超时时间。 | Duration | 否 | 30s | 无。 |
request.max-retries | 单次 OpenAPI 请求失败的最大重试次数。 | Integer | 否 | 3 | 无。 |
request.retry-after.max-delay | 可接受的 Retry-After 退避时长上限。 | Duration | 否 | 60s | 超过该值时请求按失败处理。 |
scan.discovery.rate-limit.requests-per-second | 发现类请求的进程级限流。 | Double | 否 | 5 | 同一 app-key 的所有客户端共享。 |
scan.metadata.rate-limit.requests-per-second | 节点元数据请求的进程级限流。 | Double | 否 | 20 | 同一 app-key 的所有客户端共享。 |
scan.content.rate-limit.requests-per-second | Block API 请求的进程级限流。 | Double | 否 | 10 | 同一 app-key 的所有客户端共享。 |
scan.fetch.rate-limit.requests-per-second | 全局拉取限流。 | Double | 否 | 10 | 在 Source 全部子任务之间分配。 |