本节摘要:本节装配管线第①段。
semantica/ingest/共 29 个文件、18497 行,是全项目最大的单一模块之一:本地文件、网页与 RSS、关系数据库(PostgreSQL/MySQL/SQLite/Oracle/DuckDB/MongoDB)、企业数仓(Snowflake、Databricks)、消息流(Kafka/Kinesis/RabbitMQ/Pulsar)、邮件(IMAP/POP3)、Git 仓库、MCP 资源、HuggingFace 数据集、Parquet/Arrow 等 25+ 种数据源,全部收敛到统一的数据源适配器模式之下。本节以FileIngestor为主线精读适配器三件套(类型探测/目录扫描/过滤),再看企业级摄取器的共性,最后回到 ingest 配置与编排器的协作方式。
内容来源:原项目源码 semantica/ingest/(file_ingestor.py、methods.py、registry.py、config.py)
⚠️ 注意:ingest 段产出的是「原始文件对象」(Raw Documents),不含解析后的文本——PDF 里是什么样进来的还是什么样。真正的文本抽取在 parse 段。另外 25+ 种摄取器中相当一部分依赖可选包(如 snowflake-connector、confluent-kafka),缺依赖只会在用到该摄取器时报
ConfigurationError,不影响其他数据源。
先看家底。ingest/ 目录按行数排序,头部是这些摄取器:
methods.py 1541 行 # 统一函数入口:get_ingest_method 调度 repo_ingestor.py 1138 行 # Git 仓库摄取 databricks_ingestor.py 1070 行 # Unity Catalog + Delta Lake db_ingestor.py 962 行 # 关系数据库通用摄取 snowflake_ingestor.py 923 行 # Snowflake 数仓 web_ingestor.py 912 行 # 网页 email_ingestor.py 899 行 # 邮件 IMAP/POP3 public_api_ingestor.py 881 行 # 公共 REST API feed_ingestor.py 836 行 # RSS/Atom stream_ingestor.py 665 行 # Kafka/Kinesis/RabbitMQ/Pulsar mcp_ingestor.py 452 行 # MCP 资源 huggingface_ingestor.py 322 行 # HF 数据集
设计上是教科书式的适配器模式:每个数据源一个 XxxIngestor 类,对外只暴露统一的「摄取」语义,输出的都是带元数据的 FileObject。为了让调用方不必记住 25 个类名,methods.py:1282 提供了统一函数 ingest(sources, source_type=None)——不指定类型时按字符串特征自动探测:
# methods.py:1333 起,自动探测的判定链 if source_str_lower.startswith(("http://", "https://")): if any(ext in source_str_lower for ext in [".xml", "/feed", "/rss", "/atom"]): source_type = "feed" # feed URL else: source_type = "web" # 普通网页 elif source_str_lower.startswith(("postgresql://", "mysql://", "sqlite://", ...)): source_type = "db" # 连接串 → 数据库 elif source_str_lower.startswith(("https://github.com", "https://gitlab.com")): source_type = "repo" # Git 仓库 elif source_str_lower.endswith((".ttl", ".owl", ".rdf", ".jsonld", ".n3", ".nt")): source_type = "ontology" # 本体文件 elif source_str_lower.endswith((".parquet", ".pq")): source_type = "parquet" else: source_type = "file" # 兜底:本地文件
于是 ingest("https://example.com/feed.xml") 与 ingest("doc.pdf") 是同一个函数;返回字典的顶层键随类型变化(file→"files"、web→"content"、db→"data")。更底层的 methods.py:1434 的 get_ingest_method 则先查 registry.py 里的自定义注册表,再落到内置实现,用户可以用 method_registry.register("ingest", "my_source", fn) 挂自己的数据源。
还有一个容易被忽略的安全件:ssrf.py(753 行)。web/api/feed 这类「按 URL 摄取」的入口都过它做 SSRF 防护——内网地址、云元数据端点(169.254.169.254)等在正式摄取前被拦截。企业级代码的自觉程度,往往就在这种模块里。
本地文件是最常用的入口。file_ingestor.py:81 的 FileTypeDetector 用三种策略识别文件类型,可靠性递增:
def detect_type(self, file_path, content=None) -> str: # 方法1:扩展名(最快,最常用) extension = file_path.suffix.lstrip(".").lower() if extension: return extension # 方法2:MIME 类型(文件存在时) if file_path.exists(): mime_type, _ = mimetypes.guess_type(str(file_path)) ... # 方法3:magic number 文件头(最可靠,需要内容) if content: file_type = self._detect_by_magic_numbers(content) ...
magic number 表(file_ingestor.py:190)值得一看——用文件头字节做指纹:
magic_numbers = { b"\x25\x50\x44\x46": "pdf", # %PDF b"\x50\x4b\x03\x04": "zip", # ZIP,也是 DOCX/XLSX/PPTX 的容器 b"\x89\x50\x4e\x47": "png", # PNG b"PAR1": "parquet", # Apache Parquet b"ARROW1\x00\x00": "arrow", # Apache Arrow IPC }
扩展名可以造假、MIME 会缺失,但文件头不会骗人——「.pdf 后缀实际是加密包」这类脏数据在第三道关被拦下。目录扫描则是 ingest_directory(file_ingestor.py:468):
def ingest_directory(self, directory_path, recursive=True, **filters): ... # 递归:是否扫描子目录 files = self.scan_directory(directory_path, recursive=recursive, **filters) for idx, file_info in enumerate(files, 1): try: file_obj = self.ingest_file(file_info["path"], **file_info) file_objects.append(file_obj) # 进度条带 ETA:每处理一个文件更新一次 self.progress_tracker.update_progress( tracking_id, processed=idx, total=total_files, ...) except Exception as e: if self.config.get("fail_fast", False): raise ProcessingError(...)
过滤参数通过 **filters 透传给 scan_directory(file_ingestor.py:735 附近):recursive 控制是否下钻子目录、file_types 白名单扩展名、ignore_hidden 跳过隐藏文件、大小上下限过滤异常文件。单个文件失败默认记日志继续——与编排器 build_knowledge_base 的容错策略一脉相承。
对比一下各摄取器的「重量」很有意思:文件摄取 802 行,而 Snowflake 923 行、Databricks 1070 行。重的部分不在「读数据」,在治理。Databricks 摄取器要处理 PAT/OAuth M2M 两种认证、Unity Catalog 的 catalog/schema/table 三级内省,还把 Delta Lake 的表血缘(lineage)一并取回;Snowflake 摄取器对应 warehouse/database/schema 三级加 key-pair/OAuth 认证。README 对此的定位是:让「已经躺在数仓里的表」直接变成带溯源的图节点,而不是先导出成 CSV 再导一次——省掉一跳 export/import,血缘不丢。
数据库摄取器(db_ingestor.py,962 行)则是通用关系型路线:吃 postgresql:// 风格连接串,把查询结果集变成行级记录。它和数仓摄取器的分工是——db_ingestor 管「任意库的任意查询」,snowflake/databricks 管「自家的目录、认证与血缘深度集成」。
流式摄取(stream_ingestor.py)与邮件摄取(email_ingestor.py)则代表另一类形态:数据不是「一批文件」而是「持续到达的事件」。它们的输出同样是统一形态,只是把「拉」改成了「订阅」。MCP 摄取器(mcp_ingestor.py+mcp_client.py 共约 1000 行)让 Semantica 作为 MCP 客户端从任意 MCP 服务器拉资源(仅支持 URL 形态连接,见 config 里的 MCP_SERVER_URL)——这是 0.6.x 时代很「生态化」的数据源。repo_ingestor(1138 行)支持 GitHub/GitLab 地址甚至 SCP 风格路径,把整个代码库摄取成文档流。
ingest 段在两种姿势下被使用。直连姿势是自己 new 一个摄取器:
from semantica.ingest import FileIngestor ingestor = FileIngestor(config={"fail_fast": False}) files = ingestor.ingest_directory("./docs", recursive=True, file_types=["pdf", "docx"])
配置层(ingest/config.py,201 行)的 IngestConfig 走「配置文件→环境变量→默认值」三级回退链:支持 YAML/JSON/TOML(读 ingest: 与 ingest_methods: 两段),环境变量用 INGEST_ 前缀(如 INGEST_RECURSIVE、INGEST_MAX_FILE_SIZE、INGEST_RESPECT_ROBOTS——网页摄取尊重 robots 协议也做成了配置项)。方法级 kwargs 再覆盖实例配置,优先级是「方法参数 > 实例配置 > 环境变量/文件」。这套 config.py 套路在整个项目 27 个子模块里几乎复刻了 27 遍,学会了这里是通法。
编排器姿势则由 build_knowledge_base(sources=[...]) 接管(见 1 章 02 节):_validate_sources 先过滤掉不存在的路径与非 http(s) URL(orchestrator.py:693),随后逐源调 run_pipeline,每个源都注册独立的进度跟踪与 pipeline_id。两种姿势共享同一份 config.get("ingest", {}) 配置段。
以量级收尾:一条 ingest_directory → parse → split 的微管线(支柱页的学习目标之一)只需要十几行代码,因为 ingest 输出的 FileObject 带有 content(bytes)与 text 属性(file_ingestor.py:56,UTF-8 优先、latin-1 兜底的解码),parse 段可以直接接着吃。每份摄取产物还会被 ingest_provenance.py 记下「从哪个数据源来」——溯源从管线第一段就开始积累,这是第 7 章 PROV-O 的第一块拼图。
💡 装配要点:本节装上管线的「进料口」。三个记忆点:① 适配器模式收敛 25+ 数据源,统一函数
ingest()按 task 分发,registry 支持自定义数据源;② FileIngestor 三重类型探测(扩展名→MIME→magic number)与目录过滤(recursive/file_types/ignore_hidden)是脏数据第一道闸;③ Snowflake/Databricks 摄取器最重的部分是认证、目录内省与血缘——「表原地变图节点,省一跳导出」是企业卖点。
methods.py 的 ingest() 统一入口+get_ingest_method 注册表查询;ssrf.py 给 URL 类摄取加防护。ingest_directory:recursive 递归、**filters 过滤、单文件失败默认不中断、进度带 ETA。下一节:
02 parse 解析与 normalize 归一——五类解析器如何把 PDF/CSV/代码/网页/邮件变成统一文本,以及为什么归一化必须发生在抽取之前。