ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

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

2026/9/12 16:24:52 拓冰建站 浏览量
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-zoomcampData 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 StoreRAG flow 通过{{ kv(GEMINI_API_KEY) }}从 KV Store 读取密钥所以第一步是拿到密钥并存入 KV Store。访问 Google AI Studiohttps://aistudio.google.com/app/apikey用 Google 账户登录点击 Create API Key复制生成的 key。课程文档的警告是不要把 API key 提交到 Git应使用环境变量或 Kestra 的 KV Store。在 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 UIhttp://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 dont 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 流程的三个环节可以对照文档中的过程理解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表示重新摄取清理后重建这与课程给出的实践建议定期重新摄取以保持信息最新配套。chat_with_ragio.kestra.plugin.ai.rag.ChatCompletion先按同样的 embedding provider 和 KV Store 配置检索相关内容再把检索到的上下文连同systemMessage、prompt一起交给gemini-2.5-flash生成回答。systemMessage中要求模型上下文里没有的信息要明说这是课程用来约束回答不虚构的写法。log_resultsio.kestra.plugin.core.log.Log把上一步的{{ outputs.chat_with_rag.textOutput }}打印到日志作为人工核对的输出点。注意 flow 里所有apiKey都写成{{ kv(GEMINI_API_KEY) }}而不是硬编码密钥这是文档强调的敏感信息管理方式。执行并验证输出在 Kestra UI 中打开11_chat_with_ragflow点击Execute。观察执行过程第一个任务抓取文档、创建并存储 embeddings第二个任务带着从 KV Store 检索到的上下文调用 LLM。打开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.1和postgres:188080 端口被占用时把 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),仅供参考