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。
- 访问 Google AI Studio(https://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 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 流程的三个环节,可以对照文档中的过程理解:
ingest_release_notes(io.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_rag(io.kestra.plugin.ai.rag.ChatCompletion):先按同样的 embedding provider 和 KV Store 配置检索相关内容,再把检索到的上下文连同systemMessage、prompt一起交给gemini-2.5-flash生成回答。systemMessage中要求模型"上下文里没有的信息要明说",这是课程用来约束回答不虚构的写法。log_results(io.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: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),仅供参考