news 2026/9/12 16:24:45

Data Engineering Zoomcamp 如何用 Kestra 构建基于 KV Store 的 RAG 问答流水线

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Data Engineering Zoomcamp 如何用 Kestra 构建基于 KV Store 的 RAG 问答流水线

Data Engineering Zoomcamp 如何用 Kestra 构建基于 KV Store 的 RAG 问答流水线

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

Data Engineering Zoomcamp 的第 2 周(Workflow Orchestration)在 Kestra 模块中提供了一个 RAG 练习:让大模型在回答"Kestra 1.1 发布了哪些功能"这类问题时,先从你指定的文档中检索内容,再把检索到的上下文注入 prompt,从而避免模型仅凭训练数据给出过时或错误的回答。本文的目标是在本地 Kestra 实例中跑通这条流水线:用IngestDocument任务从外部 URL 抓取发布说明、生成 embeddings 并存入 Kestra 的 KV Store,再用ChatCompletion任务带 RAG 上下文提问,最后在 Logs 标签页核对输出质量。

前提条件(课程文档明确要求):

  • Kestra 已在本地运行(kestra/kestra:v1.1镜像 +postgres:18,不要用kestra/kestra:develop,它是可能含 bug 的开发版本);
  • 一个能访问 Gemini API 的 Google 账户(文档提到有免费额度)。

准备工作:获取 Gemini API Key 并写入 KV Store

RAG flow 通过{{ kv('GEMINI_API_KEY') }}从 KV Store 读取密钥,所以第一步是拿到密钥并存入 KV Store。

  1. 访问 Google AI Studio(https://aistudio.google.com/app/apikey),用 Google 账户登录,点击 "Create API Key",复制生成的 key。课程文档的警告是:不要把 API key 提交到 Git,应使用环境变量或 Kestra 的 KV Store。
  2. 在 Kestra 中用 KV 写入任务保存密钥。仓库中 06_gcp_kv.yaml 展示了io.kestra.plugin.core.kv.Set任务的写法,用它为GEMINI_API_KEY建一个 key:
- id: set_gemini_key type: io.kestra.plugin.core.kv.Set key: GEMINI_API_KEY kvType: STRING value: <你的 Gemini API key>

其中value替换为你在第 1 步复制的密钥。这个 flow 运行一次即可,作用是让后续 RAG flow 里的kv('GEMINI_API_KEY')能取到值。

创建 RAG flow:三个任务各自做什么

在 Kestra UI(http://localhost:8080)中新建 flow,粘贴仓库提供的 11_chat_with_rag.yaml 内容。完整定义如下:

id: 11_chat_with_rag namespace: zoomcamp tasks: - id: ingest_release_notes type: io.kestra.plugin.ai.rag.IngestDocument description: Ingest Kestra 1.1 release notes to create embeddings provider: type: io.kestra.plugin.ai.provider.GoogleGemini modelName: gemini-embedding-001 apiKey: "{{ kv('GEMINI_API_KEY') }}" embeddings: type: io.kestra.plugin.ai.embeddings.KestraKVStore drop: true fromExternalURLs: - https://raw.githubusercontent.com/kestra-io/docs/refs/heads/main/src/contents/blogs/release-1-1/index.md - id: chat_with_rag type: io.kestra.plugin.ai.rag.ChatCompletion description: Query about Kestra 1.1 features with RAG context chatProvider: type: io.kestra.plugin.ai.provider.GoogleGemini modelName: gemini-2.5-flash apiKey: "{{ kv('GEMINI_API_KEY') }}" embeddingProvider: type: io.kestra.plugin.ai.provider.GoogleGemini modelName: gemini-embedding-001 apiKey: "{{ kv('GEMINI_API_KEY') }}" embeddings: type: io.kestra.plugin.ai.embeddings.KestraKVStore systemMessage: | You are a helpful assistant that answers questions about Kestra. Use the provided documentation to give accurate, specific answers. If you don't find the information in the context, say so. prompt: | Which features were released in Kestra 1.1? Please list at least 5 major features with brief descriptions. - id: log_results type: io.kestra.plugin.core.log.Log message: | ✅ RAG Response (with retrieved context): {{ outputs.chat_with_rag.textOutput }}

三个任务对应 RAG 流程的三个环节,可以对照文档中的过程理解:

  1. ingest_release_notesio.kestra.plugin.ai.rag.IngestDocument:从fromExternalURLs列出的外部 URL 抓取 Kestra 1.1 发布说明,用gemini-embedding-001模型生成 embeddings,并以io.kestra.plugin.ai.embeddings.KestraKVStore作为向量存储存入 KV Store。drop: true表示重新摄取(清理后重建),这与课程给出的实践建议"定期重新摄取以保持信息最新"配套。
  2. chat_with_ragio.kestra.plugin.ai.rag.ChatCompletion:先按同样的 embedding provider 和 KV Store 配置检索相关内容,再把检索到的上下文连同systemMessageprompt一起交给gemini-2.5-flash生成回答。systemMessage中要求模型"上下文里没有的信息要明说",这是课程用来约束回答不虚构的写法。
  3. log_resultsio.kestra.plugin.core.log.Log:把上一步的{{ outputs.chat_with_rag.textOutput }}打印到日志,作为人工核对的输出点。

注意 flow 里所有apiKey都写成{{ kv('GEMINI_API_KEY') }},而不是硬编码密钥,这是文档强调的敏感信息管理方式。

执行并验证输出

  1. 在 Kestra UI 中打开11_chat_with_ragflow,点击Execute
  2. 观察执行过程:第一个任务抓取文档、创建并存储 embeddings;第二个任务带着从 KV Store 检索到的上下文调用 LLM。
  3. 打开Logs标签页,查看log_results输出。课程文档给出的判断标准是:带 RAG 的回答应当"具体、详细、准确列出该版本真实发布的功能,并且基于实际文档"(Specific and detailed / Accurate / Grounded in actual documentation)。文档没有给出固定的期望文本,是否达标由你对照 Kestra 1.1 发布说明自行核对。

可选对照实验:仓库同时提供 10_chat_without_rag.yaml,它用普通的io.kestra.plugin.ai.completion.ChatCompletion任务提出同一个问题、不带任何检索上下文。先跑它再跑 RAG flow,可以在 Logs 里直接对比:不带上下文时,文档预期你会看到回答"含糊笼统、可能不正确、缺失具体细节",因为模型只能依赖可能过时的训练数据。这个对照就是 RAG 价值的直接演示。

限制与后续

课程在 RAG 部分给出了三条实践建议,也是这套流水线的边界说明:

  • 保持文档更新:定期重新摄取(drop: true的任务每次都会重建 embeddings)以确保信息是最新的;
  • 合理分块:大文档应切成有意义的 chunk;
  • 测试检索质量:验证被检索到的是正确的文档内容。

另外两个执行层面的注意事项:fromExternalURLs指向的是 Kestra 官方文档仓库中 1.1 发布说明的原始 Markdown URL,需要本地环境能访问该地址,摄取任务才能拿到内容;flow 的namespace固定为zoomcamp,与课程其他 flow 保持一致即可。

如果 Kestra 本身遇到问题,课程 README 的排查建议是:确认镜像固定在kestra/kestra:v1.1postgres:18;8080 端口被占用时把 Kestra 端口映射改为 18080 并访问 http://localhost:18080/;仍不工作时停止并移除现有 Kestra + Postgres 容器,用docker-compose up -d重新启动。

【免费下载链接】data-engineering-zoomcampData Engineering Zoomcamp is a free 9-week course on building production-ready data pipelines. Join the course here 👇🏼项目地址: https://gitcode.com/GitHub_Trending/da/data-engineering-zoomcamp

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/12 16:24:20

ESP32与HC-SR501人体感应模块零基础实战:从原理到感应灯

如果你第一次接触 ESP32&#xff0c;又恰好想做个“人一来就有反应”的装置——比如进门自动亮灯、感应报警器、或者统计房间里有没有人——那你大概率会搜到一块叫HC-SR501的小板子。它非常便宜&#xff0c;几块钱一块&#xff0c;圆形白色透镜&#xff0c;翠绿色的电路板&…

作者头像 李华
网站建设 2026/9/12 16:22:46

AI辅助学术写作:2025届毕业生的高效工具链

1. 2025届学术写作的AI革命2025届毕业生正面临学术写作范式的重大变革。我最近指导了几位本科生使用AI工具完成毕业论文&#xff0c;从开题到答辩仅用传统方法1/3的时间就达到了优秀水平。这并非个例——根据Nature最新调研&#xff0c;85%的顶尖高校教授已接受AI辅助的学术论文…

作者头像 李华
网站建设 2026/9/12 16:22:08

组合优化与凸优化实验指南:从梯度下降到拉格朗日对偶的实践路径

简介&#xff1a;哈工大组合优化与凸优化研究生课程实验资料包&#xff0c;面向修读该课程的研究生以及希望系统入门优化领域的算法工程师。资料以动手实验为主线&#xff0c;配合详尽的实验报告和说明书&#xff0c;帮助读者将拉格朗日对偶、KKT条件、梯度下降、牛顿法等理论概…

作者头像 李华
网站建设 2026/9/12 16:18:56

G31触发延迟导致的高速跳步测量误差与补偿实战

1. 这不是理论推演&#xff0c;是产线实测踩出来的坑G31、高速跳步、触发信号延迟、测量精度、补偿——这五个词凑在一起&#xff0c;不是实验室里的仿真波形图&#xff0c;而是我去年在某汽车零部件厂做在线尺寸检测系统升级时&#xff0c;连续三周睡在车间里盯出来的故障日志…

作者头像 李华