ADK 示例项目数据摄取管线实战:基于 KFP 与 Vertex AI Vector Search 2.0 构建 RAG 索引
ADK 示例项目数据摄取管线实战基于 KFP 与 Vertex AI Vector Search 2.0 构建 RAG 索引【免费下载链接】adk-samplesA collection of sample agents built with Agent Development Kit (ADK)项目地址: https://gitcode.com/GitHub_Trending/ad/adk-samples本文以 adk-samples 仓库中core/rag-vector-search示例项目的数据摄取模块data_ingestion/为主线系统讲解如何用 Kubeflow PipelinesKFP编排一条加载 → 分块 → 摄取的完整链路把文档导入 Agent Platform Vector SearchVector Search 2.0Collection 并自动生成 Embedding。读完本文你将掌握该管线的本地运行、远程提交、定时调度三种使用方式以及每个配置参数与底层源码组件的对应关系可直接复制命令到自己的 RAG 项目中落地。管线概览一条命令完成 RAG 索引构建在传统 RAG 方案中构建索引通常需要自行串联数据清洗 → 文本分块 → 调用 Embedding 模型 → 写入向量库多个环节。而本示例的摄取管线把整个工作流压缩为两个 KFP 组件并将Embedding 生成完全交给 Vector Search 2.0 Collection——Collection 上配置了 Embedding 模型后摄取数据对象时向量由系统自动生成无需在管线中手工调用 embedding API。从源码 pipeline.py 可以看到整条管线由两个组件按依赖顺序串联且每个组件都通过.set_retry(num_retries2)声明了最多 2 次自动重试提升了在本地与云端的稳定性# Process the data processed_data process_data(...).set_retry(num_retries2) # Ingest the processed data into Vector Search 2.0 Collection ingest_data( project_idproject_id, locationlocation, collection_idcollection_id, input_tableprocessed_data.output, ingestion_batch_sizeingestion_batch_size, ).set_retry(num_retries2)两个组件职责清晰组件源码文件职责process_datacomponents/process_data.py从 BigQuery 取数、HTML 转 Markdown、文本分块、去重结果写入 BigQuery 表ingest_datacomponents/ingest_data.py读取 BigQuery 中的分块结果批量创建 Vector Search 2.0 Collection 数据对象ingest_data组件的文档字符串点明了核心设计Embeddings are auto-generated by the Collections configured embedding model——即向量由 Collection 侧配置的 Embedding 模型示例 .env.example 中为MODEL_NAME_EMBEDDINGgemini-embedding-001在摄取时自动生成。前置条件三件准备事项原文档列出了两条前提结合仓库实际文件完整准备清单如下安装 terraform用于后面的基础设施Infra步骤即创建 Vector Search 2.0 Collection 与管线 GCS Bucket。准备.env文件在示例根目录core/rag-vector-search复制 .env.example 为.env并填入自己的项目、区域与 Collection ID。.env中与摄取管线直接相关的变量如下变量说明示例值PROJECT_IDGCP 项目 IDTODO: update-this-valueREGIONVertex AI Pipelines 区域us-central1VECTOR_SEARCH_COLLECTION_IDVector Search 2.0 Collection IDrag-vector-search-collectionVECTOR_SEARCH_LOCATIONVector Search 位置缺省回落到REGIONus-central1SERVICE_ACCOUNT远程执行使用的服务账号TODO: update-this-valuePIPELINE_ROOT管线根目录GCS 路径gs://...PIPELINE_NAME管线展示名TODO: update-this-valueCRON_SCHEDULE定时调度表达式TODO: update-this-valueDISABLE_CACHING是否禁用管线缓存falseSCHEDULE_ONLY仅创建/更新调度而不立即运行false切换到目标 GCP 项目执行gcloud config set project YOUR_PROJECT_ID确保后续命令作用在正确的项目上。从 Makefile 可以看到make>data-ingestion: set -a . ./.env set a \ cd data_ingestion uv run python data_ingestion_pipeline/submit_pipeline.py --local \ --project $$PROJECT_ID --region $$REGION --collection-id $$VECTOR_SEARCH_COLLECTION_ID快速开始两条命令打通全流程原文档明确要求两步都从示例根目录core/rag-vector-search执行。第一步provision Collectionmake setup-inframake setup-infra该命令实际执行 Makefile 中定义的动作setup-infra: cd infra/terraform terraform init terraform apply -var-filevars/env.tfvars即进入 infra/terraform 目录执行terraform init与terraform apply。它会完成三件事创建 Vector Search 2.0 CollectionTerraform 通过 scripts/setup_vector_search_collection.py 完成创建对应 vector_search.tf 等配置创建管线使用的 GCS Bucket启用所需 API见 apis.tf。执行前需要编辑 vars/env.tfvars 填入你自己的项目与区域信息。第二步运行数据摄取make>make>def run_local(args: argparse.Namespace) - None: from kfp import local local.init(runnerlocal.SubprocessRunner(use_venvFalse)) ... pipeline( project_idargs.project, locationargs.vector_search_location, schedule_timedatetime.now(UTC).isoformat(), collection_idargs.collection_id, )本地模式下 submit_pipeline.py 仅校验三个必填参数project_id、region、collection_id任一缺失都会打印错误并sys.exit(1)退出其余参数全部使用默认值。前提说明即便使用本地运行模式组件内部仍然需要访问 BigQuery 与 Vector Search 2.0因此本地环境必须配置好 Google Cloud 凭据ADC且数据摄取所依赖的 Collection 已由第一步创建完成。参数全景从pipeline.py看管线可调项管线的全部可调参数定义在 pipeline.py 的dsl.pipeline函数签名中。理解这些参数是定制自己数据摄取流程的入口参数默认值作用说明project_id无GCP 项目 IDlocation无Vector Search 区域schedule_time1970-01-01T00:00:00Z调度时间戳该默认值会被process_data识别为未设置并自动替换为当前 UTC 时间is_incrementalTrue是否仅处理最近的数据False时全量处理look_back_days1增量模式下向前回溯的天数决定处理窗口chunk_size1500文本分块大小字符数chunk_overlap20相邻分块之间的重叠字符数max_rows100最多抓取的行数0表示不限制destination_datasetrag_vector_search_qa_data结果写入的 BigQuery 数据集destination_tableincremental_questions_embeddings增量结果表deduped_tablequestions_embeddings去重结果表最终用于摄取collection_idVector Search 2.0 Collection IDingestion_batch_size250每次批量摄取请求的数据对象数量上限 250这些参数在两个组件中分别生效分块相关参数chunk_size、chunk_overlap、max_rows、增量相关作用于process_dataingestion_batch_size作用于ingest_data。其中process_data内部对默认schedule_time的处理逻辑位于 components/process_data.pyif schedule_time_dt.year 1970: logging.warning(Pipeline schedule not set. Setting schedule_time to current date.) schedule_time_dt datetime.now(timezone.utc) START_DATE schedule_time_dt - timedelta(dayslook_back_days) END_DATE schedule_time_dt也就是说无论本地运行还是云端调度处理窗口总是[schedule_time - look_back_days, schedule_time]这个闭区间。深入组件数据是如何被加工并摄入的process_data取数 → 转 Markdown → 分块 → 去重components/process_data.py 是一个基于python:3.11-slim镜像的 KFP 组件运行流程如下取数默认通过 BigQuery 内联 CTE 生成一份合成 QA 数据集12 条 Python 编程问答保证管线开箱即用、端到端跑通按is_incremental与日期窗口过滤并按max_rows截断。预处理按last_edit_date降序排序并drop_duplicates(question_id)保留每个问题的最新版本。HTML → Markdown借助markdownify将问题正文与答案转成 Markdown并把问题标题拼为 H1、每条答案拼为## Answer N:的 H2 小节最终合并为full_text_md整段文本。分块使用langchain_text_splitters的RecursiveCharacterTextSplitter(chunk_sizechunk_size, chunk_overlapchunk_overlap)对整段文本递归切分通过swifter并行 apply 加速。生成稳定 chunk_id以question_id__序号形式为每个分块生成确定性 ID使重复执行保持幂等——重跑时ingest_data会跳过已存在的对象而不是累积重复数据。源码注释也提醒若某文档分块数量减少旧分块不会被自动清理。写回 BigQuery创建按creation_timestamp做 DAY 级时间分区的表is_incrementalTrue时以append模式写入增量表随后在去重表中按question_id取creation_timestamp最新记录并以replace模式写入。输出工件将去重表地址写入 KFP 输出工件output_table供下游ingest_data读取。ingest_data批量创建数据对象Embedding 自动生成components/ingest_data.py 读取上游去重表通过google-cloud-vectorsearch的DataObjectServiceClient调用batch_create_data_objects批量创建数据对象。关键实现细节批量上限 250源码中明确注释 Max 250 per request for auto-embeddings因此batch_size min(ingestion_batch_size, 250)即使传入更大值也会被钳制到 250。向量留空请求体中vectors: {}Embedding 由 Collection 配置的模型在服务端自动生成。幂等跳过捕获google.api_core.exceptions.AlreadyExists遇到已存在的分块直接跳过并计数避免重复摄入。数据字段每个数据对象携带question_id、text_chunk、full_text_md三个字段与检索侧search_collection的output_fields遥相呼应。两个组件所需的依赖bigframes、google-cloud-vectorsearch、langchain-text-splitters、markdownify、swifter等由组件注解packages_to_install声明容器在运行时会自动安装完整清单见 data_ingestion/pyproject.toml。高级用法一提交到 Vertex AI Pipelines 远程执行当数据量大或需要云端托管时去掉--local直接调用submit_pipeline.py并补齐远程参数即可。原文档给出的命令如下cd data_ingestion uv run python data_ingestion_pipeline/submit_pipeline.py \ --project $PROJECT_ID --region $REGION \ --collection-id $VECTOR_SEARCH_COLLECTION_ID \ --service-account $SERVICE_ACCOUNT \ --pipeline-root $PIPELINE_ROOT \ --pipeline-name $PIPELINE_NAME远程模式下 submit_pipeline.py 会强制校验六个必填参数project_id、region、service_account、pipeline_root、pipeline_name、collection_id缺失即报错退出。其执行流程为用kfp.compiler.Compiler().compile()把pipeline编译为data_processing_pipeline.json构造aiplatform.PipelineJob默认enable_cachingTrue可用--disable-caching关闭通过带指数退避重试的submit_and_wait_pipeline提交并等待完成backoff.on_exception注解最多尝试 3 次、最长 1 小时重试间隔指数增长结束后删除临时编译产物。所有参数都支持从环境变量读取默认值见下表与.env中的变量一一对应命令行参数对应环境变量--projectPROJECT_ID--regionREGION--vector-search-locationVECTOR_SEARCH_LOCATION缺省回落REGION--collection-idVECTOR_SEARCH_COLLECTION_ID--service-accountSERVICE_ACCOUNT--pipeline-rootPIPELINE_ROOT--pipeline-namePIPELINE_NAME--disable-cachingDISABLE_CACHING--cron-scheduleCRON_SCHEDULE--schedule-onlySCHEDULE_ONLY完整参数列表可随时通过uv run python data_ingestion_pipeline/submit_pipeline.py --help查看。高级用法二定时调度让索引持续保鲜向量索引需要随数据更新保持新鲜原文档专门指出脚本支持--cron-schedule与--schedule-only两个调度参数。二者配合的语义如下--schedule-only只创建或更新调度不立即执行管线--cron-schedule 0 2 * * *定义 Cron 表达式如每天凌晨 2 点。调度逻辑位于 submit_pipeline.py 末尾它基于编译好的PipelineJob构造PipelineJobSchedule先按display_name查询是否已存在同名调度存在则调用update(cron...)更新表达式不存在则调用create(cron..., service_account...)新建。若--schedule-only未搭配--cron-schedule使用脚本会报 Missing --cron-schedule argument for scheduling 并退出。仓库还提供了完整的 CI/CD 示例 deployment/cloudbuild.yaml演示如何在 Cloud Build 中以SCHEDULE_ONLYTRUE的方式配合CRON_SCHEDULE创建/更新周期性的PipelineJobSchedule适合挂接到合并到 main 分支等触发器上实现自动重排程。手动触发的命令格式为gcloud builds submit --config deployment/cloudbuild.yaml \ --substitutions\ _PROD_PROJECT_IDmy-project,\ _REGIONus-central1,\ _VECTOR_SEARCH_COLLECTION_IDrag-vector-search-collection,\ _PIPELINE_GCS_ROOTgs://my-project-rag-vector-search-rag,\ _PIPELINE_SA_EMAILrag-vector-search-ragmy-project.iam.gserviceaccount.com,\ _PIPELINE_NAMErag-vector-search-ingestion,\ _CRON_SCHEDULE0 2 * * *其中_PIPELINE_GCS_ROOT引用的 Bucket 与_PIPELINE_SA_EMAIL对应的服务账号均由make setup-infra创建见 Terraform 输出的pipeline_gcs_bucket_name。验证成果测试你的 RAG 应用原文档强调管线成功完成后即可用 Vector Search 2.0 测试你的 RAG 应用。在 app/retrievers.py 中search_collection()通过DataObjectSearchServiceClient发起语义检索request vectorsearch_v1beta.SearchDataObjectsRequest( parentcollection_path, semantic_searchvectorsearch_v1beta.SemanticSearch( search_textquery, search_fieldtext_embedding, task_typeRETRIEVAL_QUERY, top_ktop_k, output_fieldsvectorsearch_v1beta.OutputFields( data_fields[question_id, text_chunk, full_text_md] ), ), )可见检索侧查询的text_embedding字段、返回的text_chunk/full_text_md字段正是摄取侧写入的数据对象字段——摄取与检索形成了完整的闭环。该函数被 app/agent.py 中的 ADK Agent 以工具形式接入因此验证方式有两种集成测试make test会运行 tests/integration/test_agent.py此时检索器被INTEGRATION_TESTTRUE环境变量替换为 Mock不访问真实 Collection但 Agent 仍会发起真实的 Gemini 调用因此需要 ADC 凭据与 Vertex AI 访问权限无凭据时测试自动跳过。交互体验make install make playground启动 ADK Web UIuv run adk web . --port 8501 --reload_agents选择app目录后即可向 Agent 提问验证基于摄取数据的回答质量。自定义数据源从示例数据到你的真实数据默认管线通过 BigQuery 内联 CTE 生成合成 QA 数据其设计意图见 process_data.py 的模块文档字符串是开箱即跑、无需外部设置。接入真实数据时只需做两处替换将fetch_sample_data中的 CTE 换成对你自己源表的查询BigQuery 支持直接查询 Cloud Storage 上的数据省去加载步骤根据源表 schema 调整下游处理逻辑HTML→Markdown 转换、分块、去重。对于大规模数据集还可以把 Markdown 转换与分块逻辑下沉为 BigQuery Python UDF利用 BigQuery 的分布式并行能力在云端执行对应参考见组件文档字符串中的官方指引。这一改造路径完全不需要改动管线编排结构pipeline.py中的参数契约保持稳定。小结围绕core/rag-vector-search示例的data_ingestion模块本文完整覆盖了从环境准备、Terraform 基建、本地/远程执行、定时调度到检索验证的全链路两条命令起步make setup-infra建 Collection 与 Bucketmake contenteditable="false">【免费下载链接】adk-samplesA collection of sample agents built with Agent Development Kit (ADK)项目地址: https://gitcode.com/GitHub_Trending/ad/adk-samples创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考