Cognee 数据源连接器实战:基于 DLT 摄取子系统的 Gmail、Slack、Notion 等 SaaS 数据接入与增量同步
【免费下载链接】cogneeCognee is the open-source AI memory platform for agents. Give your AI agents persistent long-term memory across sessions with a self-hosted knowledge graph engine.项目地址: https://gitcode.com/GitHub_Trending/co/cognee
本篇技术指南以 cognee 仓库中的 examples/integrations/README.md 为主线,系统讲解 cognee 数据源连接器(Data-source connectors)的架构定位、统一保证、Gmail 快速上手以及如何自行编写一个全新连接器。读完本文,你将掌握:连接器为何以cognee-community-connector-<source>社区包形式分发、如何用一行cognee.remember(...)把外部 SaaS 数据摄入内存知识图谱、primary_key/write_disposition/max_rows_per_table等参数如何控制增量同步与源端删除传播,以及通过DOCUMENT_SOURCE_ATTR让邮件、页面正文走标准 cognify 实体抽取的文档模式原理。
一、连接器是什么:社区包 + DLT 摄取子系统
连接器(Connector)负责把外部数据源(Gmail、Slack、Notion、Google Drive、Confluence 等)拉取进 cognee 的内存知识图谱。它们在设计上有一个关键决定:连接器以独立社区包的形式分发,而不是打进 cognee 核心。
正如 examples/integrations/README.md 所述,所有连接器都托管在topoteretes/cognee-community社区仓库下,包名为cognee-community-connector-<source>。这样做的直接好处是:cognee 核心代码不会因每个数据源而引入各自的第三方 SDK 依赖——Gmail 的 Google API、Slack 的 SDK、Notion 的客户端都不会污染核心依赖树。
这一点在核心代码中有明确的印证:cognee/tasks/ingestion/connectors/init.py 的模块文档写道:
Connectors that pull an external source (Gmail, Slack, Notion, Google Drive, Confluence, …) into cognee memory are distributed asseparate community packages… so core stays free of per-source SDKs. No connector is bundled in core.
与此同时,每一个连接器都构建在 cognee 的 DLT(dlt / data load tool)摄取子系统之上。这意味着所有连接器共享同一套摄取保证,而不是各自重新发明一套摄取逻辑。核心的调用链位于 cognee/tasks/ingestion/resolve_dlt_sources.py,其模块文档明确了这套链路的职责:resolve_dlt_sources -> ingest_dlt_source -> orphan_cleanup(见 connectors/init.py 第 9-11 行的 docstring)。
二、所有连接器共享的四大统一保证
正因为连接器统一复用 DLT 摄取子系统,用户拿到任何连接器都能获得一致的语义保证:
1. 一次调用完成摄入(One call to ingest)
连接器暴露的是一个dltsource 工厂函数,你只需把它交给cognee.remember(...)即可,无需了解任何 dlt 内部细节:
await cognee.remember(gmail_source(...), dataset_name="gmail_inbox")remember()是 cognee 的顶层内存写入入口,实现在 cognee/api/v1/remember/remember.py。对于 DLT 输入,它在内部依次执行add()(摄入)与cognify()(构建知识图谱)两步流水线。当数据是DltResource/DltSource/SourceFactory类型时(见 resolve_dlt_sources.py 的类型判定),会先经由ingest_dlt_source完成行级摄取。
2. 增量重同步(Incremental re-sync)
连接器通过write_disposition="merge"实现按主键 upsert 的增量同步,或用replace实现全量快照式同步。重复运行同一连接器时,只会拉取增量(delta),已摄入且未变化的数据不会重复处理。
在源码层面,write_disposition默认值为"replace",由remember()透传给resolve_dlt_sources(见 resolve_dlt_sources.py):
primary_key = kwargs["primary_key"] if "primary_key" in kwargs else None write_disposition = kwargs["write_disposition"] if "write_disposition" in kwargs else "replace"值得注意的是,DLT 源的data_id是稳定派生的——关系型清单源的data_id由(dataset, source name)决定,不带内容哈希(见 resolve_dlt_sources.py),因此「同一来源数据变化后重新摄入」会原地更新对应的 Data 记录,而不是产生孤儿数据。行级节点的稳定 id 则由dlt:{table}:{pk_value}:{content_hash}派生(见 _dlt_row_identifier),保证未变化的行跨运行保持同一 id、不触发重复 cognify,而编辑过的行获得新 id。
3. 源端删除传播(Forget-on-source-deletion)
如果上游删除了某条记录,连接器重新同步后,这条记录会从 cognee 的图谱(graph)+ 向量(vector)+ 关系(relational)三层存储中一并删除。这一能力由共享的orphan_cleanup路径实现。
resolve_dlt_sources返回(data, orphan_cleanup)二元组,其中orphan_cleanup是一个异步回调(见 resolve_dlt_sources.py),其核心逻辑在 _delete_dlt_orphans:
- 它只清理
system_metadata["source"]命中指定标签的 Data 记录——关系型清单默认扫("dlt", "dlt_source"),文档源则传入自己的文档标签(如("notion",)),从而互不误删; - 对清单源还会用
manifest_source_names进一步限定来源名,保证「一个数据集里多个 DLT 源,重摄其中一个不会删掉其他源」; - 删除动作被设计为延迟到新行提交之后执行,避免「先删后写」造成的数据丢失窗口(见 resolve_dlt_sources.py 的注释说明);
- 删除时通过
delete_data_nodes_and_edges同时清理图谱节点/边与向量存储,并调用delete_data删除关系型记录,还会尽力使依赖这些图谱元素的会话答案失效(见 resolve_dlt_sources.py)。
一个关键的安全边界是:当回读到的行集合为空时,孤儿清理会被跳过(do_manifest_cleanup = write_disposition != "append" and bool(manifest_data_ids)、do_document_cleanup = bool(document_fresh_ids),见 resolve_dlt_sources.py)。因为空回读无法与「同步失败/配置错误」区分,若把空集合当作「全部都是孤儿」会把整个语料库清空;保留陈旧行一轮是更安全的失败模式。
4. 散文内容以文档形式摄入(Prose ingested as documents)
邮件、Slack 消息、Notion 页面这类自然语言内容与关系型数据库行不同——它们应该通过标准 cognify 流程(LLM 实体抽取)进入图谱,而不是走关系型 schema 路径。连接器通过设置dlt_utils.DOCUMENT_SOURCE_ATTR选择文档路径,这一机制的详细原理见下文第五节。
三、可用连接器一览
从 PyPI 直接安装即可使用,无需克隆社区 monorepo:
| 数据源 | 安装包 |
|---|---|
| Gmail | cognee-community-connector-gmail |
| Slack(导出) | cognee-community-connector-slack |
| Confluence | cognee-community-connector-confluence |
| Notion | cognee-community-connector-notion |
| Google Drive | cognee-community-connector-google-drive |
每个社区包的README.md与examples/目录提供了各自数据源的安装设置、增量重同步说明以及隐私/选择加入(opt-in)注意事项。
四、Gmail 快速上手
以 Gmail 连接器为例,完整的摄入 + 检索流程如下:
pip install cognee-community-connector-gmailimport cognee from cognee_community_connector_gmail import gmail_source await cognee.remember( gmail_source(label_ids=["INBOX"], credentials_path="credentials.json"), dataset_name="gmail_inbox", primary_key="id", write_disposition="merge", max_rows_per_table=0, # 0 = no read cap, so forget-on-delete sees the whole inbox ) answer = await cognee.search( query_text="What did my manager ask me to do this week?", datasets=["gmail_inbox"], )各参数含义如下:
gmail_source(...):连接器暴露的 dlt source 工厂函数。label_ids=["INBOX"]限定要拉取的邮件标签,credentials_path指向 OAuth 凭据文件。注意:gmail 连接器是懒加载第三方 SDK 的,因此首次 import 不会拉起 Google API 依赖。dataset_name="gmail_inbox":目标数据集。数据集不存在时会自动创建(见 remember.py 的resolve_authorized_user_datasets调用)。后续可用datasets=["gmail_inbox"]限定检索范围。primary_key="id":增量同步与去重的行主键。write_disposition="merge"时按此键 upsert。默认值为"id"(见 resolve_dlt_sources.py 的primary_key or "id"回退逻辑)。write_disposition="merge":dlt 写入策略,merge表示按主键合并增量;replace表示每次全量重写(适合没有删除流的快照型源,如 Notion/Slack 导出)。第三个取值是append——追加模式会跳过清单源的孤儿清理(见 resolve_dlt_sources.py),因为每次运行本就预期新增行。max_rows_per_table=0:单表最大读取行数上限,0表示不设上限。正如文档注释强调的,设置为 0 才能让 forget-on-delete 看到整个收件箱——如果截断了行数,超出上限的旧行会被误判为「不在当前语料中」从而被清理,或反之导致已删除的邮件无法被识别为孤儿。cognee.search(...):摄入完成后即可对知识图谱做语义检索,返回的answer即基于已建图的回答。
此外,remember()还支持run_in_background=True(后台异步运行,返回的RememberResult可await等待完成)、chunk_size/chunker(分块控制)、custom_prompt(自定义实体抽取提示词)等参数,详见 remember.py 的参数文档。
五、文档模式深入:DOCUMENT_SOURCE_ATTR的原理
这是连接器机制中最值得深入的一层。核心常量定义在 cognee/tasks/ingestion/dlt_utils.py:
DOCUMENT_SOURCE_ATTR = "cognee_document_source" def document_source_tag(item) -> Optional[str]: tag = getattr(item, DOCUMENT_SOURCE_ATTR, None) return tag if isinstance(tag, str) and tag else None机制:一个 dlt source 通过给自身设置cognee_document_source属性来声明「我的每一行是一条文本文档」;属性值就是该源行数据system_metadata["source"]的标签(例如"notion"、"google_drive")。这让resolve_dlt_sources保持连接器无关——连接器自行声明自身性质,核心引擎不硬编码任何连接器名称。
在 resolve_dlt_sources.py 中,dlt 源被分成两组:
- 文档源(document-tagged):每个行变成一个文本文档 DataItem,走标准 cognify(LLM 实体抽取)。行被期望携带
title/content列(可选url/id),构造时title会被转成# 标题前缀拼到正文前面,url写入system_metadata["url"]、id写入external_id(见 _build_document_data_item)。 - 关系型源(relational):整个源被折叠成一个清单(manifest)DataItem,是一个描述该源全部唯一行的 JSON 文档,后续由专门的 DLT cognify 流水线(
extract_dlt_source_edges)从关系 schema 确定性构造图谱边,而不是走 LLM。
为什么要有这个区分?核心注释(dlt_utils.py)解释得很清楚:文档模式让每一行成为一段文本,流入正常的 cognify 实体抽取——这正是邮件、消息、页面这类自然语言内容想要的;而关系型数据的图谱是确定性地从 schema 构建的(跳过 LLM,避免为每行关系数据花 token)。
单元测试 cognee/tests/unit/tasks/test_dlt_document_mode.py 精确验证了这个接缝:
- 没有标记的源
document_source_tag返回None,走关系型路径;空字符串/非字符串标记同样被忽略; - 标记为
"notion"/"google_drive"的行,is_dlt_sourced返回False(因为其source != "dlt"),从而被classify_documents路由到TextDocument/cognify,而不是清单 schema 路径; _build_document_data_item正确产出带# My Page标题前缀、携带url/external_id元数据的文档项。
此外还有集成测试 cognee/tests/integration/tasks/test_dlt_orphan_graph_vector_purge.py 与 test_add_foreground_orphan_cleanup.py 验证孤儿数据在图谱/向量/关系三层的真实清除行为。
一个值得注意的文档模式设计:文档源尊重调用方的write_disposition,两种同步模型都可用——replace用于没有删除流的快照源(Notion/Slack,每次运行用当前可见行重写暂存区),merge(配合hard_delete墓碑列)用于有真实删除流的增量源(如 Google Drive 的 Changes API)。无论哪种方式,都会把整个集合读回(max_rows_per_table=0),这样孤儿清理才能识别掉出当前语料库的行(见 resolve_dlt_sources.py 的注释)。
六、编写一个新的连接器
参考现有连接器作为模板,发布一个cognee-community-connector-<source>包。连接器需要满足以下契约:
# cognee_community_connector_<source>/__init__.py from dlt.sources import DltSource def <source>_source(**config) -> DltSource: """返回一个 dlt source,声明 primary_key / write_disposition / hard_delete。""" ...必备要素:
- 工厂函数返回一个
dltsource。它声明:primary_key:行的稳定主键(增量 upsert 与去重依据);write_disposition:merge(有增量语义)或replace(全量快照);hard_delete标记列:用于指示删除的墓碑列——配合merge让 delete feed 驱动的源能够把上游删除传播给 cognee 的孤儿清理路径。
- 散文类源设置
DOCUMENT_SOURCE_ATTR:如果源的内容是邮件、消息、页面等自然语言,就给 source 设置cognee_document_source属性(值为你的源标签,如"gmail"),使行以文档形式摄入、走 LLM 实体抽取,而不是关系型 schema 路径。注意行需要携带title/content列(可选url/id),见 _build_document_data_item 对列名的约定。 - 第三方 SDK 保持懒导入:在函数内部 import 各自的 SDK,避免导入包本身时就拉起 Gmail/Notion/Slack 等重量级依赖。
- 配套 mocked 测试:带上「mocked SaaS + mocked LLM」的测试——CI 中不使用真实凭据。仓库核心侧的 test_dlt_document_mode.py 用
SimpleNamespace模拟源对象、无数据库无 LLM 地验证文档模式接缝,是编写这类测试的绝佳范本。
写完后在 PyPI 发布cognee-community-connector-<source>包,并在cognee-community仓库中附带该包的README.md与examples/,说明安装设置、增量重同步与隐私/选择加入注意事项。
七、小结
cognee 的数据源连接器体系是一个「小而美」的架构设计:核心只提供一套 DLT 摄取引擎(resolve_dlt_sources -> ingest_dlt_source -> orphan_cleanup),每个外部数据源以社区包形式接入并自行声明摄取语义。统一保证(一次调用摄入、增量重同步、源端删除传播、散文文档化)让使用者无需关心每个数据源的接入细节,而DOCUMENT_SOURCE_ATTR机制则让「关系型数据走确定性 schema 建图、自然语言走 LLM 实体抽取」两种路径在同一个引擎下优雅共存。
如果你想深入了解相关实现,可以继续阅读仓库中的这些位置:
- 连接器包约定:cognee/tasks/ingestion/connectors/init.py
- DLT 摄取核心链路:cognee/tasks/ingestion/resolve_dlt_sources.py
- 文档模式常量与工具:cognee/tasks/ingestion/dlt_utils.py
remember()入口与参数路由:cognee/api/v1/remember/remember.py- 文档模式单元测试:cognee/tests/unit/tasks/test_dlt_document_mode.py
- 孤儿清理集成测试:cognee/tests/integration/tasks/test_dlt_orphan_graph_vector_purge.py
【免费下载链接】cogneeCognee is the open-source AI memory platform for agents. Give your AI agents persistent long-term memory across sessions with a self-hosted knowledge graph engine.项目地址: https://gitcode.com/GitHub_Trending/co/cognee
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考