钉钉知识库 CDC(公测中)

更新时间:
复制 MD 格式

本文为您介绍如何使用钉钉知识库 CDC 连接器。

背景信息

钉钉知识库 CDC 连接器通过周期性轮询钉钉开放平台 OpenAPI,按文档修改时间持续捕获钉钉知识库(Wiki Workspace)中在线文档的新增和更新。连接器输出文档的标题、链接、所属知识库、父节点、更新时间等元数据和 Markdown 格式正文,也可以选择获取文档的原始块结构,适用于 RAG 知识库构建、文档检索索引和湖仓内容同步等场景。

钉钉知识库没有可供直接消费的变更事件流,钉钉知识库 CDC 通过 OpenAPI 轮询实现全量和增量读取:

  1. 每个知识库是一个独立的读取分片。未配置 workspace-id 时,连接器按 scan.discovery.interval 周期调用知识库列表接口发现操作人可见的全部知识库;配置后仅同步指定知识库。

  2. 每次拉取时,连接器按 scan.discovery.max-depthscan.discovery.max-nodes 遍历知识库节点树,筛选出钉钉在线文档节点(type=FILEcategory=ALIDOC),并按文档修改时间生成本次拉取的游标位置。

  3. 游标之后的文档按修改时间排序,单个分片单次最多输出 scan.prefix-size 篇。对每篇变化的文档,连接器通过一致性校验机制读取其块结构并渲染成 Markdown 正文。

  4. 无增量变化时,连接器按 scan.poll-interval 间隔重复上述检查。

连接器按主键 doc_id 输出 upsert 语义的 changelog,包含 INSERT 和 UPDATE_AFTER,不输出 UPDATE_BEFORE:

场景

输出

全量阶段首次读到的文档

INSERT

增量阶段内容或元数据变化的文档

UPDATE_AFTER

当前实现不感知文档删除:文档被删除后不再出现在节点列表中,连接器不会输出 DELETE 事件,下游中该行会保留为最后一次同步的状态。如需清理残留数据,可提交批模式作业,以下游表与源表按主键做反连接(anti-join),删除主键已不存在于源表中的行。

钉钉 CDC 连接器支持的信息如下。

类别

详情

支持类型

源表、数据摄入数据源

运行模式

流模式;数据摄入 YAML 另支持批模式(仅 snapshot 启动模式)

数据格式

固定 Schema,不支持裁剪字段

API 种类

SQL 和数据摄入 YAML

同步对象

知识库中操作人有权读取的钉钉在线文档节点(节点类型 type=FILEcategory=ALIDOC

启动模式

initial(全量加增量)、snapshot(仅有界全量)、latest-offset(仅增量);不支持 timestamp

Watermark

不支持

维表和结果表

不支持

前提条件

网络要求

Flink 作业需要访问钉钉开放平台 OpenAPI(默认域名为 api.dingtalk.com)。作业部署在 VPC 内时,请确保具备公网访问能力(参见阿里云实时计算 Flink 版文档网络连接)。

应用创建与接口授权

无论使用哪种启动模式,都需要先在钉钉开放平台完成应用创建和接口授权。

  1. 登录钉钉开发者后台,创建企业内部应用,在应用的"凭证与基础信息"页面获取应用的 Client ID(AppKey)和 Client Secret(AppSecret)。连接器使用 AppKey 和 AppSecret 调用 POST /v1.0/oauth2/accessToken 获取企业内部应用 access_token 并自动续期。

  2. 为应用开通本连接器所需的接口权限。在开发者后台进入目标应用的"权限管理"页面,搜索下表中的权限点并申请开通(部分权限点需要组织管理员审批)。连接器会调用以下钉钉 OpenAPI:

API

用途

权限点

POST /v1.0/oauth2/accessToken

获取企业内部应用 access_token

基础能力,无需额外权限点

GET /v2.0/wiki/workspaces

列出操作人可见的知识库

知识库读权限

GET /v2.0/wiki/nodes?parentNodeId=

遍历知识库节点树,列出子节点

知识库节点读权限

GET /v2.0/wiki/nodes/{nodeId}

查询节点元数据

知识库节点读权限

GET /v1.0/doc/suites/documents/{docKey}/blocks

查询文档块元素

企业存储文件读权限

权限点名称以钉钉开放平台官方文档为准,汇总如下:

能力

对应 API

所需权限

是否必需

获取访问凭证

oauth2/accessToken

无需额外权限点

读取知识库列表

wiki/v2.0/wiki/workspaces

知识库读权限

读取节点列表和节点元数据

wiki/v2.0/wiki/nodes

知识库节点读权限

读取文档块内容

doc/v1.0/doc/suites/documents/{docKey}/blocks

企业存储文件读权限

操作人与知识库数据授权

operator-id 填写一名企业成员的 Union ID(可通过钉钉开放平台"查询用户详情"接口获取),连接器所有接口都以该操作人身份调用,可见范围由操作人自身的权限决定:

  • 操作人必须对目标知识库或文档至少拥有"仅可查看"权限("可查看/下载"及以上亦可);

  • 若作业未配置 workspace-id,连接器只能发现操作人可见的知识库。

钉钉知识库支持知识库级、文件夹级和单文档级三层权限,文档默认继承上级目录权限。权限类型及本连接器的最低要求如下:

权限类型

能否被本连接器读取

说明

仅可查看

满足最小读取要求

可查看/下载

推荐,避免个别接口对下载类能力的额外要求

可编辑

超出读取所需,不建议仅为同步授予

可管理

超出读取所需,不建议仅为同步授予

无权限

不能

该知识库或文档不会出现在同步结果中

授权方式:由知识库管理员在知识库的"设置 > 成员及权限"中添加操作人(支持按人、部门、群授权),或在单个文档的分享设置中为操作人开放权限。新建文档默认继承知识库权限,按整库授权后新增文档无需单独配置。连接器为只读链路,不需要任何写权限。

授权问题排查

连接器对授权类错误不会重试(会直接失败或跳过对应文档),请按下表排查:

错误码 / 现象

含义

处理方式

forbidden.acrossOrg

跨组织访问被拒绝

确认应用、操作人和目标知识库属于同一组织;跨组织知识库无法同步

orgAuthLevelNotEnough

组织认证等级不满足接口要求

在钉钉管理后台完成企业认证或提升认证等级后重试

权限点不足(permission denied 类报错)

应用未开通所需权限点

回到"权限管理"页面核对知识库读权限、知识库节点读权限、企业存储文件读权限是否均已开通

知识库 / 文档不出现在同步结果中

操作人无权读取

在知识库"成员及权限"中为操作人添加"仅可查看"及以上权限

nodeNotExist(404)

节点已删除或 ID 失效

连接器自动跳过,无需处理

doc key is illegal

文档标识非法

属于永久失败,不会重试,请核对节点状态或联系钉钉支持

使用限制

  • 仅支持源表,不支持维表和结果表,不支持 Watermark。

  • 只同步知识库中的钉钉在线文档节点(type=FILEcategory=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-intervalscan.prefix-size 和限流参数控制调用成本与新内容的同步延迟。

  • 进程级限流参数(scan.discovery/metadata/content.rate-limit.requests-per-second)由同一 app-key 的所有客户端共享,多个作业共用同一应用时请整体评估配额。

  • content 是从文档块结构渲染出的 Markdown 文本,不保证完整保留所有富媒体和结构化元素(无法渲染的块会以占位注释形式标注);需要原始块结构时,Flink SQL 声明 blocks_json METADATA 列,数据摄入 YAML 开启 scan.block.include-raw-json 输出 blocks_json 物理列。

  • 故障恢复可能重复输出 INSERT 和 UPDATE_AFTER,建议下游按主键 doc_id 幂等写入。

  • 资源发现是累加的:知识库或文档一旦被作业发现,即使之后调小 workspace-id 或追加 scan.document.exclude-node-ids,也不会从已有作业中移除。需要彻底移除时,请停止作业并无状态重新启动。

  • 授权类错误(forbidden.acrossOrgorgAuthLevelNotEnough、非法 docKey 等)属于永久失败,连接器不会自动重试,请参见"授权问题排查"处理。

SQL

特色功能

  • 固定 Schema 源表:输出文档标题、链接、正文(Markdown)和元数据,按主键 doc_id 输出 upsert 语义的 changelog。

  • 启动模式:通过 scan.startup.mode 选择读取方式。

模式

是否读取全量

是否读取增量

首次启动行为

initial

先全量读取范围内的文档,再持续轮询增量变化

snapshot

只全量读取范围内的文档,完成后结束

latest-offset

只处理作业启动后发生的文档变化

  • 原始块结构输出:通过 blocks_json METADATA 列输出文档原始块结构的 JSON 文本,需同时开启 scan.block.include-raw-json

语法结构

Flink SQL DDL 必须声明下表中全部 7 个物理字段,字段顺序、名称和大小写需保持一致(字段名大小写敏感),并声明 PRIMARY KEY (doc_id) NOT ENFORCED

DDL 字段

类型

说明

doc_id

STRING NOT NULL

文档节点标识,源表主键

workspace_id

STRING

文档所属知识库 ID

parent_node_id

STRING

文档在知识库节点树中的父节点 ID

title

STRING

文档标题

url

STRING

文档访问链接

content

STRING

文档正文,由文档块结构渲染成的 Markdown 文本

modified_time

TIMESTAMP_LTZ(3)

文档最后修改时间

此外,DDL 可以通过 METADATA 声明以下只读元数据。元数据列不改变上述物理 Schema。

Metadata Key

类型

说明

blocks_json

STRING

文档原始块结构的 JSON 文本,可用于保留块级结构或自定义解析。声明该列时必须同时设置 'scan.block.include-raw-json' = 'true',二者必须同时启用或同时缺省

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

固定值为 dingtalk-cdc

object-type

同步的对象类型。

String

当前仅支持 document

app-key

钉钉企业内部应用的 Client ID(AppKey),用于获取 access_token。

String

无。

app-secret

钉钉企业内部应用的 Client Secret(AppSecret)。

String

无。

operator-id

操作人的 Union ID,该用户必须有权读取待同步的文档。

String

可通过钉钉开放平台"查询用户详情"接口获取。

endpoint

钉钉开放平台 OpenAPI 域名。

String

https://api.dingtalk.com

无。

workspace-id

将发现范围限定到指定知识库。

String

支持逗号分隔多个值。

说明

不配置时自动发现操作人可见的全部知识库。建议显示配置workspace-id

scan.document.exclude-node-ids

排除的节点 ID 列表。

String

逗号分隔,这些节点在发现和拉取阶段都会被跳过。

scan.discovery.max-depth

单次知识库轮询扫描允许的最大节点树深度。

Integer

100

无。

scan.discovery.max-nodes

单次知识库轮询扫描允许访问的最大节点数。

Integer

100000

无。

scan.startup.mode

启动模式。

String

initial

支持 initialsnapshotlatest-offset,不支持 timestamp

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

是否通过 blocks_json 只读元数据列输出原始块 JSON。

Boolean

false

声明 blocks_json METADATA 列时必须设为 true。

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=FILEcategory=ALIDOC),其他对象会被跳过;不感知文档删除。

语法结构

YAML 固定 Schema 如下,主键为 doc_id。与 Flink SQL 相比,YAML Schema 中 blocks_json 作为物理列直接输出(而非 METADATA 列),且 workspace_id 为 NOT NULL。

字段

类型

说明

doc_id

STRING NOT NULL

文档节点标识,主键

workspace_id

STRING NOT NULL

文档所属知识库 ID

parent_node_id

STRING

文档在知识库节点树中的父节点 ID

title

STRING

文档标题

url

STRING

文档访问链接

content

STRING

文档正文,由文档块结构渲染成的 Markdown 文本

blocks_json

STRING

文档原始块结构的 JSON 文本,仅当 scan.block.include-raw-json: true 时输出

modified_time

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

固定值为 dingtalk-cdc

object-type

同步的对象类型。

String

当前仅支持 document

app-key

钉钉企业内部应用的 Client ID(AppKey)。

String

无。

app-secret

钉钉企业内部应用的 Client Secret(AppSecret)。

String

无。

operator-id

操作人的 Union ID,该用户必须有权读取待同步的文档。

String

无。

endpoint

钉钉开放平台 OpenAPI 域名。

String

https://api.dingtalk.com

无。

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

支持 initialsnapshotlatest-offset,不支持 timestamp;批模式仅支持 snapshot

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 中包含 blocks_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 全部子任务之间分配。