基于 ADK 的混合多格式 RAG 数据摄入流水线:从 GCS 到 Vertex AI Vector Search 2.0 的自动化实践
【免费下载链接】adk-samplesA collection of sample agents built with Agent Development Kit (ADK)项目地址: https://gitcode.com/GitHub_Trending/ad/adk-samples
本文聚焦 contrib/python/multiformat-hybrid-rag/data_ingestion_pipeline/README.md 所讲解的 Data Ingestion Pipeline:一个将异构格式文档自动化摄入 Vertex AI Vector Search 2.0(VS2)的 KFP(Kubeflow Pipelines)流水线。文章以该 README 为骨架,结合仓库内 pipeline.py、submit_pipeline.py、Makefile 及
src/下核心实现逐层展开,读者读完可掌握:如何搭建 VS2 基础设施、如何提交与调度摄入流水线、三个处理阶段(preprocess → chunk_and_index → cleanup)的内部原理,以及全量参数与配置项的实操含义。
流水线定位:为多格式混合 RAG 自动构建向量索引
在 multiformat-hybrid-rag 这个 ADK 示例应用中,检索增强生成(RAG)的质量高度依赖索引数据的时效性与完整性。Data Ingestion Pipeline 正是负责这一环的自动化引擎:它把散落在 GCS 桶中的 PDF、Office 文档、Markdown、HTML、JSON 等异构文件,自动完成加载 → 切块(chunking)→ 导入 VS2 Collection的完整链路,而向量嵌入(embedding)由 Collection 配置的 embedding 模型自动生成,无需额外维护嵌入服务。
该流水线既可手动触发完成全量初始加载,也可通过 Cron 调度周期性运行,使搜索索引始终与源数据保持同步。从仓库的 pyproject.toml 可以看到,它的依赖集合(kfp>=2.0.0、google-cloud-pipeline-components>=2.19.0、google-cloud-vectorsearch>=0.4.0、langchain-text-splitters>=1.1.2、backoff>=2.2.0等)清晰表明这是一个构建在 Vertex AI Pipelines 之上的 KFP 工程,Python 版本要求>=3.11, <=3.13。
前置条件:项目 ID 与开发环境准备
1. 设置环境变量
所有make命令都依赖 Google Cloud Project ID,因此第一步是将其导出为环境变量:
export GOOGLE_CLOUD_PROJECT="YOUR_PROJECT_ID"将"YOUR_PROJECT_ID"替换为你的真实 GCP 项目 ID。从 Makefile 的TF_VARS定义可以看出,项目内几乎所有基础设施变量都从.env文件加载,GOOGLE_CLOUD_PROJECT是整个链路的地基;此外还会用到GOOGLE_CLOUD_LOCATION(默认us-central1)与GOOGLE_CLOUD_LOCATION_MODELS(默认global,Gemini 3.x 系列模型只在 global 端点发布)。
2. 供给基础设施(Datastore)
在仓库根目录执行:
make setup-datastore该命令负责供给 VS2 Collection、GCS 存储桶与 service account,要求本机已安装并配置好terraform。
需要说明的是,当前仓库 Makefile 中实际提供的基础设施目标是setup-infra(以及配套的tf-plan、tf-destroy)。setup-infra的执行过程可以印证 README 中"供给 VS2 Collection、GCS 桶、service account"的描述,其内部依次完成:
- 校验
.env存在,并设置 gcloud 项目与 quota project; - 启用
cloudresourcemanager、serviceusage等引导 API; - 通过
terraform init && terraform apply在 infra/terraform/dev 目录供给 API、BigQuery、Vector Search、Cloud Run 与 IAM 等全部资源; - 预创建 Cloud Build 暂存桶(避免三个并行构建竞争自动建桶而失败);
- 并行构建三个镜像(preprocess 服务、chunk-index 服务、流水线镜像)并部署两个 Cloud Run 服务。
该目标要求gcloud、terraform与uv均已安装,且.env中已填好GCS_BUCKET、SERVICE_ACCOUNT、PIPELINE_ROOT等键值。
运行数据摄入流水线
基础设施就绪后,即可运行流水线完成数据摄入。
a. 提交流水线到 Vertex AI Pipelines
在仓库根目录执行(确保当前 shell 仍保留前文导出的GOOGLE_CLOUD_PROJECT):
make>PYTHONPATH=.:data_ingestion_pipeline PIPELINE_IMAGE=$PIPELINE_IMAGE \ uv run python data_ingestion_pipeline/data_ingestion_pipeline/submit_pipeline.py \ --service-account=$SERVICE_ACCOUNT \ --pipeline-root=$PIPELINE_ROOT \ --disable-caching \ --cron-schedule="${PIPELINE_CRON:-0 2 * * *}"其中--cron-schedule默认值为0 2 * * *,即每天凌晨 2 点自动执行一次增量摄入;--disable-caching用于强制本次运行不使用 Vertex AI 的 pipeline caching。
b. 监控流水线进度
提交时 SDK 会把流水线的控制台 URL 打印到标准输出。需要详细监控时,可在 Google Cloud Console 中打开Vertex AI Pipelines面板,查看每个组件的执行状态、日志与产物。每个组件都配置了set_retry(num_retries=2)(见 pipeline.py),瞬时故障会自动重试最多 2 次;而提交层本身还套了一层backoff.expo指数退避(最多尝试 3 次、总时长上限 1 小时),可从容应对瞬时 API 错误。
流水线全参数详解:来自源码的单一口径
流水线的每个参数都以.env默认值为准,同时允许在 CLI 上按需覆盖——pipeline.py的形参、submit_pipeline.py的 argparse 参数与.env键三者保持单一事实来源(single source of truth),同一套配置同时驱动本地开发、CI 与生产运行。
下表整理自 pipeline.py 与 submit_pipeline.py:
| 参数 | 默认值 | 说明 |
|---|---|---|
--project-id | .env的GOOGLE_CLOUD_PROJECT | GCP 项目 ID,必填 |
--region | .env的GOOGLE_CLOUD_LOCATION | GCP 区域,须与 BQ dataset 与 VS2 collection 所在区域一致,必填 |
--gcs-prefix | documents/(GCS_PREFIX) | 桶内待摄入文件的目录前缀 |
--bq-dataset | rag_pipeline(BQ_DATASET) | 承载流水线各张表的 BigQuery dataset |
--vs-collection-id | multiformat-hybrid-rag-collection | VS2 中保存 chunk 的 collection ID |
--vs-documents-collection-id | multiformat-hybrid-rag-documents | VS2 中按 file_id 保存文档元数据的 collection ID |
--chunk-size | 800(CHUNK_SIZE) | 单个 chunk 的最大字符数 |
--chunk-overlap | 50(CHUNK_OVERLAP) | 相邻 chunk 之间的重叠字符数,保证边界上下文不丢失 |
--vs-batch-size | 250(VS_BATCH_SIZE) | 每次 VS2 批量创建调用的数据对象数,API 上限为 250 |
--rechunk-all | False | 强制对全部文件重新切块(如修改chunk_size/chunk_overlap后使用) |
--skip-cleanup | False | 跳过删除检测步骤 |
--service-account | .env的SERVICE_ACCOUNT | 流水线容器运行身份,须具备 BQ、GCS、VS2、Cloud Functions 权限(由 Terraform 供给),必填 |
--pipeline-root | .env的PIPELINE_ROOT | GCS 路径,Vertex AI 存放中间产物(组件输出、executor 日志等),必填 |
--pipeline-name | rag-ingestion(PIPELINE_NAME) | 流水线展示名 |
--disable-caching | False | 关闭 Vertex AI pipeline caching |
--cron-schedule | .env的CRON_SCHEDULE | 周期性执行的 Cron 表达式 |
--schedule-only | False | 只创建/更新调度,不立即执行 |
其中 BQ 表名在 pipeline.py 中通过字符串插值拼为全限定名:{project_id}.{dataset}.{gcs_objects|preprocessed|chunks},并作为参数传递给下游组件。submit_pipeline.py在启动前会对四个必填项(project、region、service account、pipeline root)做校验,缺失即报错退出并提示在.env或 CLI 中补齐。
除--disable-caching、--rechunk-all、--skip-cleanup、--schedule-only四个开关外,其余参数既可以从.env读取,也可以直接作为 CLI 参数传入。直接运行脚本的等价形式(来自 submit_pipeline.py 的用法注释):
PYTHONPATH=.:data_ingestion_pipeline uv run python \ data_ingestion_pipeline/data_ingestion_pipeline/submit_pipeline.py \ --service-account=SA_EMAIL \ --pipeline-root=gs://bucket/pipeline-root流水线内部结构:三阶段 DAG 源码拆解
data_ingestion_pipeline/pipeline.py 定义了一个三阶段有向无环图(DAG),KFP 将其编译为 JSON 规范后提交给 Vertex AI Pipelines 执行:
preprocess ──► chunk_and_index ──► cleanup (conditional)三个阶段分别对应三个 KFP 组件(@component装饰,运行在 Artifact Registry 中的同一容器镜像data-pipeline:latest内),阶段间通过.after()强制串行,各自带 2 次自动重试。值得注意的是,组件函数内的 import 语句必须写在函数体内部——KFP 只序列化函数体为临时脚本在容器内执行,编译机上的 import 不会生效。
阶段一:Preprocess —— 变更检测与文本抽取
components/preprocess.py 是薄封装,真正逻辑位于 src/document_preprocessing/preprocess.py。其核心流程:
- 变更检测:以 BQ Object Table(GCS 桶元数据镜像的外部表)LEFT JOIN
preprocessed表,找出新增文件(prep.file_id IS NULL)与内容变化文件(obj.md5_hash != prep.content_hash); - 内容去重:两层去重——跨批次去重(同一 md5 已抽取过则复用其 file_id,写入
duplicate_of:<id>的 stub 行)与批内去重(同批次多个 URI 共享 md5 时只抽取字典序最小的 URI); - 并行抽取:通过
ThreadPoolExecutor以最多 200 个 worker(PREPROCESS_MAX_WORKERS)向 preprocess Cloud Run 服务发起带认证的 HTTP fanout 请求,由服务内部用 LibreOffice + Gemini 完成文本抽取与相关性判定; - 流式落库:抽取结果按每批 100 行(
PREPROCESS_FLUSH_BATCH_SIZE)流式写入带 run_id 后缀的 staging 表,最后以单条 MERGE(按 file_id upsert)合并进preprocessed表。
file_id = MD5(gcs_uri)是文件的确定性身份标识——它在 Python(hashlib)与 BigQuery(TO_HEX(MD5(uri)))两侧可一致计算,且与文件内容无关,文件更新后 ID 保持不变。
该阶段支持的文件类型包括:PDF、DOCX、DOC、PPTX、PPT、XLSX、XLS、RTF、HTML、JSON、JSONL、Markdown 与纯文本;解析器实现见 src/document_preprocessing/parser/ 目录。变更检测 SQL 中的扩展名白名单与解析器的PARSEABLE_MIMES必须保持同步,否则未被收录的格式会被静默跳过。
阶段二:Chunk & Index —— 切块、上下文摘要与向量化
components/chunk_and_index.py 是整个流水线计算量最大的步骤,其逻辑实现在 src/chunking/chunk_and_index.py:
- 候选识别:以"纯元数据扫描 + 按 file_id 聚类回表取内容"的两阶段查询找出需要(重新)切块的文件——从未索引过、或
extracted_at > 最近一次成功索引时间的文件。这里比较的是indexed_at(仅在 VS2 确认成功后盖章)而非chunked_at,因此 VS2 写入失败的文件天然会在下一轮被重试; - Markdown 感知切块:chunk-index Cloud Run 服务使用
langchain-text-splitters进行切块——尊重#/##/###标题作为自然边界,超长段落依次回退到按段落、按行、按词切分,相邻 chunk 间保留chunk_overlap字符的重叠以保证上下文连续; - 上下文摘要生成:每个 chunk 通过 Gemini 生成一段"相对全文而言该 chunk 讲了什么"的上下文摘要(contextual summary),提升检索命中质量;
- 写入 VS2 与 BQ:chunk 以空向量批量创建数据对象(
vs_batch_size上限 250),VS2 使用 collection 配置的 embedding 模型自动生成嵌入;同时 chunk 元数据(chunk_id = {file_id}__{chunk_index})写入 BQ chunks 表用于追踪调试。
编排器在服务确认新 chunk 就绪之后、才会删除该文件的旧 chunk(批量删除,CHUNK_INDEX_DELETE_BATCH_SIZE默认 100),从而保证切块服务故障时旧索引依然可用——"过期结果优于空结果"。
阶段三:Cleanup —— 删除传播(条件执行)
components/cleanup.py 负责摄入的逆操作:当源 GCS 桶中的文件被删除时,检测"孤儿"文件并级联删除全部派生数据。实现位于 src/removal/propagate_gcs_deletions.py。
删除顺序对幂等性至关重要:先删 VS2 数据对象 → 再删 BQ chunks 表 → 最后删 BQ preprocessed 表。这样即使步骤中途失败被重试,已删除的资源也只是 no-op,不会留下悬空数据。
该步骤被包在dsl.Condition(skip_cleanup == False, name="run-cleanup")中(pipeline.py),传入--skip-cleanup即可跳过。源码注释特别提醒:条件必须写== False而非is False,因为 KFP 会把条件编译为仅支持==比较的 CEL 表达式。
周期性调度:让索引始终保持新鲜
README 明确指出流水线可以"定期运行以确保搜索索引保持最新",具体由--cron-schedule与--schedule-only两个参数支撑(submit_pipeline.py):
- 传入
--cron-schedule="0 2 * * *"时,脚本会先立即执行一次,再创建(或更新)同名的PipelineJobSchedule;默认情况下make contenteditable="false">【免费下载链接】adk-samplesA collection of sample agents built with Agent Development Kit (ADK)项目地址: https://gitcode.com/GitHub_Trending/ad/adk-samples
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考