本文介绍如何使用飞书知识库 CDC 连接器。
背景信息
飞书 CDC 连接器通过周期性轮询飞书开放平台 OpenAPI,按文档编辑时间持续捕获飞书知识空间(Wiki Space)中在线文档的新增和更新。连接器输出文档的标题、链接、更新时间等元数据和 Markdown 格式正文,也可以获取文档的原始块结构,适用于 RAG 知识库构建、文档检索索引和湖仓内容同步等场景。
飞书 CDC 连接器支持的信息如下。
类别 | 详情 |
连接器类型 | 源表 |
API | Flink SQL、数据摄入 YAML |
同步对象 | 知识空间中应用有权读取的 |
全增量模式 |
|
增量模式 |
|
全量模式 |
|
数据格式 | 固定 Schema,不支持裁剪字段 |
Watermark | 不支持 |
维表和结果表 | 不支持 |
工作原理
飞书知识库没有可供直接消费的变更事件流,飞书 CDC 通过 OpenAPI 轮询实现全量和增量读取:
每个知识空间是一个独立的读取分片。连接器按
scan.discovery.interval周期发现新授权的知识空间。每次拉取时,连接器列出空间内的全部节点,筛选出配置范围内的
doc、docx文档节点,并按文档编辑时间生成本次拉取的游标位置。游标之后的文档按编辑时间排序,单个分片单次最多输出
scan.prefix-size篇。对每篇变化的文档,连接器读取其块结构并渲染成 Markdown 正文。无增量变化时,连接器按
scan.poll-interval间隔重复上述检查。
变更到 changelog 的映射
连接器按主键 node_id 输出 upsert 语义的 changelog,包含 INSERT 和 UPDATE_AFTER,不输出 UPDATE_BEFORE。
场景 | 输出 |
全量阶段首次读到的文档 | INSERT |
增量阶段内容或元数据变化的文档 | UPDATE_AFTER |
当前实现不感知文档删除:文档被删除后不再出现在节点列表中,连接器不会输出 DELETE 事件,下游中该行会保留为最后一次同步的状态。
已删除文档的清理任务建议独立进行,直接调用飞书 OpenAPI 完成比对,低频、定期执行,可选择周末或夜间等业务低峰时段,避免与同步作业的轮询叠加争抢 OpenAPI 限流。比对方式如下:
完整拉取一次知识空间节点列表(
GET /open-apis/wiki/v2/spaces/{space_id}/nodes,须全部分页成功),与下游已存储的 node_id 比对,找出「下游有、列表无」的节点。对这类节点逐个调用
GET /open-apis/wiki/v2/spaces/get_node二次确认:返回节点不存在类错误的,判定为真实删除;返回权限拒绝类错误的,说明文档仍存在但应用暂时失去访问权限,不得删除;限流、网络等瞬时错误则跳过本轮,下一周期再确认。确认为真实删除后执行清理,建议软删除,行数据保留用于审计和恢复;对误删零容忍的场景,可连续多个周期均确认删除后才执行清理。
前提条件
Flink 作业需要访问飞书开放平台 OpenAPI(默认域名为 open.feishu.cn)。作业部署在 VPC 内时,请确保具备公网访问能力。
无论使用哪种启动模式,都需要先在飞书开放平台完成应用创建和知识空间授权。
登录飞书开放平台,创建企业自建应用,获取应用的 App ID 和 App Secret。连接器使用 App ID 和 App Secret 获取并自动续期 tenant_access_token。
为应用开通本连接器所需的接口权限,并发布应用版本。连接器会调用以下飞书 OpenAPI:
API
用途
GET /open-apis/wiki/v2/spaces列出应用可访问的知识空间
GET /open-apis/wiki/v2/spaces/{space_id}/nodes列出知识空间中的文档节点
GET /open-apis/wiki/v2/spaces/get_node解析文档节点元数据
GET /open-apis/docx/v1/documents/{document_token}/blocks读取文档块内容
将应用添加为目标知识空间的成员,确保应用有权读取待同步的文档。具体操作参见飞书开放平台的知识库常见问题中“如何给应用授权访问知识库文档资源”。连接器只同步应用已被授权的
doc、docx节点。
权限点列表
上述接口对应的具体权限点(scope)名称以飞书开放平台文档为准。授权时需要覆盖以下两类能力:
能力 | 对应 API | 是否必需 |
读取知识空间和节点列表 |
| 是 |
读取文档块内容 |
| 是 |
使用限制
仅支持源表,不支持维表和结果表,不支持 Watermark。
只同步知识空间中应用有权读取的
doc、docx类型节点;其他对象(表格、多维表、文件等)会被跳过。当前实现不感知文档删除,不输出 DELETE 事件。
Flink SQL DDL 必须声明固定的 6 个物理字段和
PRIMARY KEY (node_id) NOT ENFORCED,不支持增删物理列;blocks_json只能以METADATA列声明。Flink SQL 一张源表只读取一个知识空间;
wiki.space-id只能填写一个值。数据摄入 YAML 支持逗号分隔多个空间或不配置自动发现。增量链路基于 OpenAPI 轮询,会产生持续的 API 调用。空间内文档较多时,请结合
scan.poll-interval、scan.prefix-size和限流参数控制调用成本与新内容的同步延迟。content是从文档块结构渲染出的 Markdown 文本,不保证完整保留所有富媒体和结构化元素;需要原始块结构时,Flink SQL 使用blocks_jsonMETADATA 列,数据摄入 YAML 直接读取blocks_json物理列。资源发现是累加的,调小
wiki.space-id、wiki.document-url或wiki.node-token不会移除已发现的知识空间或文档。故障恢复可能重复输出 INSERT 和 UPDATE_AFTER,建议下游按主键
node_id幂等写入。YAML Schema 与 SQL Schema 不同:YAML 额外包含
space_id物理列,且blocks_json为物理列而非 METADATA 列。请根据接入方式选择对应的 Schema 定义。
启动模式
通过 scan.startup.mode 选择读取方式。
模式 | 是否读取全量 | 是否读取增量 | 首次启动行为 |
| 是 | 是 | 先全量读取范围内的文档,再持续轮询增量变化 |
| 是 | 否 | 只全量读取范围内的文档,完成后结束 |
| 否 | 是 | 读取更新时间不早于 |
| 否 | 是 | 只处理作业启动后发生的文档变化 |
运行过程中新发现的知识空间或文档按 scan.new-split.startup.mode 决定读取方式:initial 先读取当前完整内容再跟随后续变化;latest-offset 只处理发现之后发生的变化。
同步范围
Flink SQL 和数据摄入 YAML 对同步范围的配置能力不同:
能力 | Flink SQL | 数据摄入 YAML |
| 必填,只能填写一个值 | 可选,支持逗号分隔多个值;不配置时自动发现全部可见空间 |
| 最多配置一个,且只能填写一篇文档 | 支持逗号分隔多个值 |
多空间同步方式 | 每个空间创建一张源表 | 单个作业内按空间自动分表( |
未配置文档选择器时,空间内应用有权读取的全部 doc、docx 文档都会被同步。
注意:数据摄入 YAML 配置文档选择器后,默认整体同步配置空间并自动发现包含所选文档的未配置空间。需要严格边界时请将 wiki.restrict-documents-to-selected-spaces 设为 true,详见「数据摄入 > 空间到表的路由」。
Flink SQL 如需把范围收窄到空间内的某一篇文档,可以二选一配置:
wiki.document-url:填写飞书文档链接,连接器从链接中解析文档标识,支持常见的wiki、docx、docs等文档地址形式。wiki.node-token:直接填写知识空间中的节点 Token。
两个参数最多配置一个,且所选文档必须属于 wiki.space-id 配置的空间。配置后,全量和增量都只处理这一篇文档。YAML 侧文档选择器支持多篇文档,行为详见数据摄入章节。
资源发现是累加的:知识空间或文档一旦被作业发现,即使之后调小范围配置,也不会从已有作业中移除。按整个空间同步的作业,空间内后续新增的文档会被正常轮询发现,不需要修改配置。需要彻底移除时,请停止作业并无状态重新启动。
SQL
语法结构
CREATE TEMPORARY TABLE <yourTableName> (
node_id STRING NOT NULL,
title STRING,
url STRING,
`type` STRING,
content STRING,
updated_time TIMESTAMP_LTZ(3),
PRIMARY KEY (node_id) NOT ENFORCED
) WITH (
'connector' = 'feishu-cdc',
'object-type' = 'document',
'app-id' = '<yourAppId>',
'app-secret' = '<yourAppSecret>',
'wiki.space-id' = '<yourWikiSpaceId>'
);Flink SQL DDL 必须声明下表中全部 6 个物理字段,字段顺序和名称需保持一致。
DDL 字段 | 类型 | 说明 |
| STRING NOT NULL | 文档节点标识,源表主键 |
| STRING | 文档标题 |
| STRING | 文档访问链接 |
| STRING | 飞书对象类型,例如 |
| STRING | 文档正文,由文档块结构渲染成的 Markdown 文本 |
| TIMESTAMP_LTZ(3) | 文档最后更新时间 |
源表必须声明 PRIMARY KEY (node_id) NOT ENFORCED。
此外,DDL 可以通过 METADATA 声明以下可读元数据。元数据列不改变上述物理 Schema。
Metadata Key | 类型 | 说明 |
| STRING | 文档原始块结构的 JSON 文本,可用于保留块级结构或自定义解析 |
WITH 参数
基础参数
参数 | 类型 | 默认值 | 是否必填 | 说明 |
| STRING | 无 | 是 | 固定为 |
| STRING | 无 | 是 | 同步的对象类型,当前仅支持 |
| STRING | 无 | 是 | 飞书应用的 App ID,用于获取 tenant_access_token |
| STRING | 无 | 是 | 飞书应用的 App Secret |
| STRING |
| 否 | 飞书开放平台 OpenAPI 域名 |
读取范围参数
参数 | 类型 | 默认值 | 是否必填 | 说明 |
| STRING | 无 | 是 | 要读取的单个知识空间 ID |
| STRING | 无 | 否 | 文档链接,将范围收窄到单篇文档;与 |
| STRING | 无 | 否 | 文档节点 Token,将范围收窄到单篇文档;与 |
启动和轮询参数
参数 | 类型 | 默认值 | 说明 |
| STRING |
| 启动模式,支持 |
| BIGINT | 无 |
|
| STRING |
| 新发现知识空间的读取方式,支持 |
| DURATION |
| 无增量变化时的轮询间隔 |
| DURATION |
| 新授权知识空间的发现间隔 |
| INTEGER |
| 单个分片单次拉取最多输出的变化文档数 |
| INTEGER | 无 | Source 并行度;不配置时使用作业默认并行度 |
请求和限流参数
参数 | 类型 | 默认值 | 说明 |
| INTEGER |
| 知识空间 OpenAPI 分页大小,范围 1~50 |
| DURATION |
| 单次 OpenAPI 请求超时时间 |
| INTEGER |
| 瞬时请求失败的最大重试次数 |
| DOUBLE |
| 列出知识空间和节点类 OpenAPI 每秒最大请求数 |
| DOUBLE |
| 解析节点元数据类 OpenAPI 每秒最大请求数 |
| DOUBLE |
| 读取文档块内容类 OpenAPI 每秒最大请求数 |
代码示例
全量加增量同步
示例将指定知识空间的文档写入 Paimon 表,声明了 blocks_json 元数据列。
CREATE TEMPORARY TABLE feishu_wiki_source (
node_id STRING NOT NULL,
title STRING,
url STRING,
`type` STRING,
content STRING,
updated_time TIMESTAMP_LTZ(3),
blocks_json STRING METADATA VIRTUAL,
PRIMARY KEY (node_id) NOT ENFORCED
) WITH (
'connector' = 'feishu-cdc',
'object-type' = 'document',
'app-id' = '<yourAppId>',
'app-secret' = '<yourAppSecret>',
'wiki.space-id' = '<yourWikiSpaceId>',
'scan.startup.mode' = 'initial'
);
CREATE TABLE IF NOT EXISTS paimon.feishu_wiki.documents (
node_id STRING,
title STRING,
url STRING,
`type` STRING,
content STRING,
blocks_json STRING,
updated_time TIMESTAMP_LTZ(3),
PRIMARY KEY (node_id) NOT ENFORCED
);
INSERT INTO paimon.feishu_wiki.documents
SELECT node_id, title, url, `type`, content, blocks_json, updated_time
FROM feishu_wiki_source;SNAPSHOT 全量扫描
SNAPSHOT 只做一次全量读取,适合内容迁移或一次性导出。
CREATE TEMPORARY TABLE feishu_wiki_snapshot (
node_id STRING NOT NULL,
title STRING,
url STRING,
`type` STRING,
content STRING,
updated_time TIMESTAMP_LTZ(3),
PRIMARY KEY (node_id) NOT ENFORCED
) WITH (
'connector' = 'feishu-cdc',
'object-type' = 'document',
'app-id' = '<yourAppId>',
'app-secret' = '<yourAppSecret>',
'wiki.space-id' = '<yourWikiSpaceId>',
'scan.startup.mode' = 'snapshot'
);如需只同步空间内的一篇文档,在 WITH 参数中追加 'wiki.document-url' = '<文档链接>' 或 'wiki.node-token' = '<节点Token>'。
数据摄入
数据摄入连接器的 type 固定为 feishu-cdc。每个知识空间对应一张名为 <知识空间ID>.documents 的表,连接器在首次发现空间时自动创建表并同步 Schema。
语法结构
source:
type: feishu-cdc
name: <yourSourceName>
app-id: <yourAppId>
app-secret: <yourAppSecret>
wiki.space-id: <yourWikiSpaceId>
scan.startup.mode: initial配置项
基础参数
参数 | 类型 | 默认值 | 是否必填 | 说明 |
| STRING | 无 | 是 | 固定为 |
| STRING | 无 | 是 | 飞书应用的 App ID,用于获取 tenant_access_token |
| STRING | 无 | 是 | 飞书应用的 App Secret |
| STRING |
| 否 | 飞书开放平台 OpenAPI 域名 |
| STRING |
| 否 | 同步的对象类型,当前仅支持 |
同步范围参数
参数 | 类型 | 默认值 | 是否必填 | 说明 |
| STRING | 无 | 否 | 逗号分隔的知识空间 ID 列表;不配置时自动发现应用有权读取的全部知识空间 |
| STRING | 无 | 否 | 逗号分隔的文档链接列表,将同步范围收窄到指定文档 |
| STRING | 无 | 否 | 逗号分隔的节点 Token 列表,将同步范围收窄到指定文档 |
| BOOLEAN |
| 否 | 是否将显式选择的文档严格限制在
|
启动和轮询参数
参数 | 类型 | 默认值 | 说明 |
| STRING |
| 启动模式,支持 |
| BIGINT | 无 |
|
| STRING |
| 新发现知识空间的读取方式,支持 |
| DURATION |
| 无增量变化时的轮询间隔 |
| DURATION |
| 新知识空间的发现间隔 |
| INTEGER |
| 单个分片单次拉取最多输出的变化文档数 |
请求和限流参数
参数 | 类型 | 默认值 | 说明 |
| INTEGER |
| 知识空间 OpenAPI 分页大小,范围 1~50 |
| DURATION |
| 单次 OpenAPI 请求超时时间 |
| INTEGER |
| 瞬时请求失败的最大重试次数 |
| DOUBLE |
| 列出知识空间和节点类 OpenAPI 每秒最大请求数 |
| DOUBLE |
| 解析节点元数据类 OpenAPI 每秒最大请求数 |
| DOUBLE |
| 读取文档块内容类 OpenAPI 每秒最大请求数 |
固定 Schema
字段 | 类型 | 说明 |
| STRING NOT NULL | 文档节点标识,主键 |
| STRING NOT NULL | 文档所属知识空间 ID |
| STRING | 文档标题 |
| STRING | 文档访问链接 |
| STRING | 飞书对象类型,例如 |
| STRING | 文档正文,由文档块结构渲染成的 Markdown 文本 |
| STRING | 文档原始块结构的 JSON 文本 |
| TIMESTAMP_LTZ(3) | 文档最后更新时间 |
固定 Schema 的主键为 node_id。与 Flink SQL 相比,YAML Schema 增加了 space_id 物理列,blocks_json 也作为物理列直接输出。
空间到表的路由
未配置 wiki.space-id 时,连接器自动发现应用有权读取的全部知识空间,每个空间生成对应的 <知识空间ID>.documents 表。配置 wiki.space-id(支持逗号分隔多个值)后,只同步指定的空间。
配置文档选择器 wiki.document-url 或 wiki.node-token(均支持逗号分隔多个值)后,选择器与 wiki.space-id 默认取并集:配置的空间仍整体同步,选择器不做收窄;未配置的空间如果包含所选文档,也会被自动发现并同步这些文档。wiki.restrict-documents-to-selected-spaces 控制是否把选择器严格限制在 wiki.space-id 配置的空间内。
假设 wiki.space-id 配置为空间 A,文档选择器包含属于 A 的文档 X 和属于 B 的文档 Y:
wiki.restrict-documents-to-selected-spaces | 同步结果 |
| 空间 A 整体同步;空间 B 被自动发现并同步文档 Y |
| 仅同步文档 X;文档 Y 不属于配置空间,作业启动时报错 |
如想让 wiki.space-id 作为严格同步边界,建议设置为 true,避免选择器静默突破边界、同步到预期外的空间。
代码示例
以下示例把指定知识空间中的文档全量读取后持续同步增量变化,写入 Paimon。
source:
type: feishu-cdc
name: Feishu Wiki Source
app-id: <yourAppId>
app-secret: <yourAppSecret>
wiki.space-id: <yourWikiSpaceId>
scan.startup.mode: initial
route:
- source-table: <yourWikiSpaceId>.documents
sink-table: feishu_wiki.documents
sink:
type: paimon
catalog.properties.metastore: rest
catalog.properties.uri: <yourDlfCatalogEndpoint>
catalog.properties.warehouse: <yourWarehouse>