基于 ADK 的混合多格式 RAG 数据摄入流水线:从 GCS 到 Vertex AI Vector Search 2.0 的自动化实践 📅 发布时间:2026/9/16 21:12:05 👁 浏览次数: 基于 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.0VS2的 KFPKubeflow 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 可以看到它的依赖集合kfp2.0.0、google-cloud-pipeline-components2.19.0、google-cloud-vectorsearch0.4.0、langchain-text-splitters1.1.2、backoff2.2.0等清晰表明这是一个构建在 Vertex AI Pipelines 之上的 KFP 工程Python 版本要求3.11, 3.13。前置条件项目 ID 与开发环境准备1. 设置环境变量所有make命令都依赖 Google Cloud Project ID因此第一步是将其导出为环境变量export GOOGLE_CLOUD_PROJECTYOUR_PROJECT_ID将YOUR_PROJECT_ID替换为你的真实 GCP 项目 ID。从 Makefile 的TF_VARS定义可以看出项目内几乎所有基础设施变量都从.env文件加载GOOGLE_CLOUD_PROJECT是整个链路的地基此外还会用到GOOGLE_CLOUD_LOCATION默认us-central1与GOOGLE_CLOUD_LOCATION_MODELS默认globalGemini 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_PROJECTmake>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_retries2)见 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_PROJECTGCP 项目 ID必填--region.env的GOOGLE_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.env的SERVICE_ACCOUNT流水线容器运行身份须具备 BQ、GCS、VS2、Cloud Functions 权限由 Terraform 供给必填--pipeline-root.env的PIPELINE_ROOTGCS 路径Vertex AI 存放中间产物组件输出、executor 日志等必填--pipeline-namerag-ingestionPIPELINE_NAME流水线展示名--disable-cachingFalse关闭 Vertex AI pipeline caching--cron-schedule.env的CRON_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-accountSA_EMAIL \ --pipeline-rootgs://bucket/pipeline-root流水线内部结构三阶段 DAG 源码拆解data_ingestion_pipeline/pipeline.py 定义了一个三阶段有向无环图DAGKFP 将其编译为 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 TableGCS 桶元数据镜像的外部表LEFT JOINpreprocessed表找出新增文件prep.file_id IS NULL与内容变化文件obj.md5_hash ! prep.content_hash内容去重两层去重——跨批次去重同一 md5 已抽取过则复用其 file_id写入duplicate_of:id的 stub 行与批内去重同批次多个 URI 共享 md5 时只抽取字典序最小的 URI并行抽取通过ThreadPoolExecutor以最多 200 个 workerPREPROCESS_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)是文件的确定性身份标识——它在 Pythonhashlib与 BigQueryTO_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 与 BQchunk 以空向量批量创建数据对象vs_batch_size上限 250VS2 使用 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, namerun-cleanup)中pipeline.py传入--skip-cleanup即可跳过。源码注释特别提醒条件必须写 False而非is False因为 KFP 会把条件编译为仅支持比较的 CEL 表达式。周期性调度让索引始终保持新鲜README 明确指出流水线可以定期运行以确保搜索索引保持最新具体由--cron-schedule与--schedule-only两个参数支撑submit_pipeline.py传入--cron-schedule0 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),仅供参考