基于 ADK 的混合多格式 RAG 数据摄入流水线:从 GCS 到 Vertex AI Vector Search 2.0 的自动化实践
2026/9/16 21:11:47 网站建设 项目流程

基于 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.0google-cloud-pipeline-components>=2.19.0google-cloud-vectorsearch>=0.4.0langchain-text-splitters>=1.1.2backoff>=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-plantf-destroy)。setup-infra的执行过程可以印证 README 中"供给 VS2 Collection、GCS 桶、service account"的描述,其内部依次完成:

  • 校验.env存在,并设置 gcloud 项目与 quota project;
  • 启用cloudresourcemanagerserviceusage等引导 API;
  • 通过terraform init && terraform apply在 infra/terraform/dev 目录供给 API、BigQuery、Vector Search、Cloud Run 与 IAM 等全部资源;
  • 预创建 Cloud Build 暂存桶(避免三个并行构建竞争自动建桶而失败);
  • 并行构建三个镜像(preprocess 服务、chunk-index 服务、流水线镜像)并部署两个 Cloud Run 服务。

该目标要求gcloudterraformuv均已安装,且.env中已填好GCS_BUCKETSERVICE_ACCOUNTPIPELINE_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.envGOOGLE_CLOUD_PROJECTGCP 项目 ID,必填
--region.envGOOGLE_CLOUD_LOCATIONGCP 区域,须与 BQ dataset 与 VS2 collection 所在区域一致,必填
--gcs-prefixdocuments/GCS_PREFIX桶内待摄入文件的目录前缀
--bq-datasetrag_pipelineBQ_DATASET承载流水线各张表的 BigQuery dataset
--vs-collection-idmultiformat-hybrid-rag-collectionVS2 中保存 chunk 的 collection ID
--vs-documents-collection-idmultiformat-hybrid-rag-documentsVS2 中按 file_id 保存文档元数据的 collection ID
--chunk-size800CHUNK_SIZE单个 chunk 的最大字符数
--chunk-overlap50CHUNK_OVERLAP相邻 chunk 之间的重叠字符数,保证边界上下文不丢失
--vs-batch-size250VS_BATCH_SIZE每次 VS2 批量创建调用的数据对象数,API 上限为 250
--rechunk-allFalse强制对全部文件重新切块(如修改chunk_size/chunk_overlap后使用)
--skip-cleanupFalse跳过删除检测步骤
--service-account.envSERVICE_ACCOUNT流水线容器运行身份,须具备 BQ、GCS、VS2、Cloud Functions 权限(由 Terraform 供给),必填
--pipeline-root.envPIPELINE_ROOTGCS 路径,Vertex AI 存放中间产物(组件输出、executor 日志等),必填
--pipeline-namerag-ingestionPIPELINE_NAME流水线展示名
--disable-cachingFalse关闭 Vertex AI pipeline caching
--cron-schedule.envCRON_SCHEDULE周期性执行的 Cron 表达式
--schedule-onlyFalse只创建/更新调度,不立即执行

其中 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。其核心流程:

  1. 变更检测:以 BQ Object Table(GCS 桶元数据镜像的外部表)LEFT JOINpreprocessed表,找出新增文件(prep.file_id IS NULL)与内容变化文件(obj.md5_hash != prep.content_hash);
  2. 内容去重:两层去重——跨批次去重(同一 md5 已抽取过则复用其 file_id,写入duplicate_of:<id>的 stub 行)与批内去重(同批次多个 URI 共享 md5 时只抽取字典序最小的 URI);
  3. 并行抽取:通过ThreadPoolExecutor以最多 200 个 worker(PREPROCESS_MAX_WORKERS)向 preprocess Cloud Run 服务发起带认证的 HTTP fanout 请求,由服务内部用 LibreOffice + Gemini 完成文本抽取与相关性判定;
  4. 流式落库:抽取结果按每批 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:

  1. 候选识别:以"纯元数据扫描 + 按 file_id 聚类回表取内容"的两阶段查询找出需要(重新)切块的文件——从未索引过、或extracted_at > 最近一次成功索引时间的文件。这里比较的是indexed_at(仅在 VS2 确认成功后盖章)而非chunked_at,因此 VS2 写入失败的文件天然会在下一轮被重试;
  2. Markdown 感知切块:chunk-index Cloud Run 服务使用langchain-text-splitters进行切块——尊重#/##/###标题作为自然边界,超长段落依次回退到按段落、按行、按词切分,相邻 chunk 间保留chunk_overlap字符的重叠以保证上下文连续;
  3. 上下文摘要生成:每个 chunk 通过 Gemini 生成一段"相对全文而言该 chunk 讲了什么"的上下文摘要(contextual summary),提升检索命中质量;
  4. 写入 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询