进阶 dagster.io 2026-10-10 01:48:56 · 6 阅读

第18章 利用 Dagster 构建低成本 AI 管道

利用 Dagster OpenAI 集成,在严格控制成本的同时释放 LLM 的潜力。

你好。随着 LLM 日益成熟且触手可及,许多组织正致力于将其融入常规业务流程,例如技术支持机器人、在线实时辅助以及各类知识库相关任务。

虽然 LLM 的理念令人惊叹且用例丰富,但运行此类流程的成本也日益凸显。

本文将分享一种方法:通过在 Dagster 新的 OpenAI 集成之上构建 AI 管道,在利用 LLM 强大能力的同时有效把控成本。

希望本教程能激励你构建创新的 AI 驱动流程,而无需挥霍预算。

一年多前,我们分享了使用 Dagster、GPT-3 和 LangChain 创建支持机器人的初步步骤。如今,随着大型语言模型(LLM)不断涌现和进化,我们开启了新篇章,融入了 GPT-4 的最新更新以及 Dagster 特别是新的 dagster-openai 集成。

如果你想跟随项目进度,可以通过我们的 GitHub 仓库进行实操。

初始设置回顾

我们的旅程始于利用 Dagster 简化 GitHub 文档的获取与索引,从而使支持机器人能精准响应用户查询。

前置条件

这不是入门教程,我们假设你熟悉 Python、命令行操作,并对 Dagster 有基本了解。

运行所需环境: - 本地安装 Python 3 - 一个 OpenAI API 密钥 - 一组 Python 依赖项,安装方式如下:

``` ### 安装 Dagster 依赖 pip install dagster dagster-aws dagster-cloud dagster-openai ### 安装 LangChain 依赖 pip install faiss-cpu langchain langchain-community langchain-openai ```

我们的初始设置从从 GitHub 仓库提取原始文档开始: @asset def source_docs(): return list(get_github_docs("dagster-io", "dagster"))

然后,将这些文档及其嵌入分块并构建为 Faiss 索引: @asset def search_index(source_docs): source_chunks = [] splitter = CharacterTextSplitter(separator=" ", chunk_size=1024, chunk_overlap=0) for source in source_docs: for chunk in splitter.split_text(source.page_content): source_chunks.append(Document(page_content=chunk, metadata=source.metadata)) search_index = FAISS.from_documents(source_chunks, OpenAIEmbeddings()) with open("search_index.pickle", "wb") as f: pickle.dump(search_index.serialize_to_bytes(), f)

最后,我们使用配合 LangChain 的向量空间搜索引擎来提升支持机器人的效率: def print_answer(question): with open("search_index.pickle", "rb") as f: serialized_search_index = pickle.load(f)

search_index = FAISS.deserialize_from_bytes(serialized_search_index, OpenAIEmbeddings()) print( chain( { "input_documents": search_index.similarity_search(question, k=4), "question": question, }, return_only_outputs=True, )["output_text"] )

该基础设施为稳健的问答机制奠定了基础。它利用 OpenAI 模型基于我们的文档解读用户问题,并结合 LangChain 构建低成本提示词。

AI 管道面临的新挑战

随着系统扩展及新模型和 AI 能力的引入,深入理解管道变得至关重要,尤其是在优化 OpenAI 服务使用方面。我们旨在制定更具成本效益的策略,重点包括:

- 增强 OpenAI 使用可见性:了解运行支持机器人及其他 AI 驱动功能的成本影响是我们的首要关切。在开发复杂 AI 管道的同时监控不同模型和 API 的使用情况,凸显了更高效追踪系统的需求。我们寻求一种方案,以便更便捷、精准地管理 OpenAI 交互。 - 成本控制:随着我们将 LLM 用于生成补全和嵌入等更复杂的任务,成本上升的担忧也随之而来。有效控制和预测成本的机制变得不可或缺。

因此,我们需要一种直接的方法,基于实际数据比较 OpenAI 模型,从而做出既能优化性能又能确保成本效益的明智决策。

介绍 dagster-openai 集成

在优化 AI 管道的过程中,我们很高兴地分享,通过开源的 dagster-openai 集成,我们将与 OpenAI API 的工作成果向更广泛的社区开放。结合 Dagster 的软件定义资产,该集成不仅促进了与 OpenAI API 的无缝交互,还通过在资产元数据中自动记录 OpenAI 使用元数据来增强透明度。欲了解详细洞察,请探索我们关于利用 Dagster 结合 OpenAI 的指南。

@asset(compute_kind="OpenAI") def openai_asset(context: AssetExecutionContext, openai: OpenAIResource): with openai.get_client(context) as client: client.chat.completions.create( model="gpt-3.5-turbo", messages=[{"role": "user", "content": "Say this is a test."}], )

该集成的突出特性是通过 Dagster Insights 放大对 OpenAI API 使用的可见性。通过系统记录使用元数据,我们现在能够对 OpenAI 模型参与情况进行更细致的分析。这不仅有助于成本优化,还使我们在评估和比较模型时能够做出更基于数据的决策。

该集成还引入了名为 with_usage_metadata 的方法,旨在记录来自任何 OpenAI 端点的使用数据。我们稍后将深入细节。

利用新集成改进方法

利用新集成改进方法,使我们能够显著提升项目的性能和成本效益:

- 可见性与控制:利用该集成,我们在 Dagster 生态系统中监控 OpenAI 使用情况,获得关于运营实践的详细洞察和警报。 - 高效性能比较:通过利用记录的元数据,我们高效地比较了模型性能,促进了决策过程。

让我们更新原有的代码!如前文初始设置所述,我们已有: @asset def search_index(source_docs): source_chunks = [] splitter = CharacterTextSplitter(separator=" ", chunk_size=1024, chunk_overlap=0) for source in source_docs: for chunk in splitter.split_text(source.page_content): source_chunks.append(Document(page_content=chunk, metadata=source.metadata)) search_index = FAISS.from_documents(source_chunks, OpenAIEmbeddings()) with open("search_index.pickle", "wb") as f: pickle.dump(search_index.serialize_to_bytes(), f)

使用新集成,只需几行代码变更: 不再直接调用 API,而是使用来自 OpenAI 资源的客户端,它会自动将 API 使用数据记录到资产目录中。 @asset(compute_kind="OpenAI") def search_index(context: AssetExecutionContext, openai: OpenAIResource, source_docs: List[Any]): source_chunks = [] splitter = CharacterTextSplitter(separator=" ", chunk_size=1024, chunk_overlap=0) for source in source_docs: for chunk in splitter.split_text(source.page_content): source_chunks.append(Document(page_content=chunk, metadata=source.metadata))

with openai.get_client(context) as client: search_index = FAISS.from_documents( source_chunks, OpenAIEmbeddings(client=client.embeddings) )

return search_index.serialize_to_bytes()

此外,我们更直观地标记了“计算类型”(compute kind)。这个小而关键的标签提升了资产图的可读性,使一眼就能理解每个资产中涉及的计算背景:

现在运行该资产。 资产物化立即添加了元数据:

更进一步,当使用 Dagster+ 运行此管道时,我们会得到漂亮的聚合图表:

接下来,更新 print_answer。首先,让我们使用 OpenAIResource 结合 LangChain,将 print answer 函数更新为补全资产: @asset(compute_kind="OpenAI") def completion(context: AssetExecutionContext, openai: OpenAIResource, search_index: Any): question = "What can I use Dagster for?" search_index = FAISS.deserialize_from_bytes(search_index, OpenAIEmbeddings()) with openai.get_client(context) as client: chain = load_qa_with_sources_chain(OpenAI(client=client.completions, temperature=0)) context.log.info( chain( { "input_documents": search_index.similarity_search(question, k=4), "question": question, }, return_only_outputs=True, )["output_text"] )

现在,你可以看到所有步骤都显示在资产图中:

只需几处代码修改,我们就实现了一个自动记录使用情况的管道,其使用数据会写入资产元数据。这为长期追踪和优化支持机器人的成本与性能奠定了坚实基础。

让我们尝试简单问题“什么是 Dagster?”:

它像以前一样工作,但提供了更多使用洞察。

现在,进入有趣的部分。

通过增强文档解决复杂查询

目前,我们的基础支持机器人在面对“如何将 dbt 与 Dagster 配合使用?”这类具体问题时,回答相当笼统:

为了增强对复杂查询(如 dagster-dbt 集成等特定主题)的响应,我们可以纳入更广泛的文档。例如,可以查阅相关的指南、API 文档和 GitHub 讨论。

为此,我们需要修改代码以从各种来源获取文档。 def get_github_docs(repo_owner, repo_name, category, archive_name="master"): with tempfile.TemporaryDirectory() as d: # 存档名称可以是分支、标签或提交。 r = requests.get(f"https://github.com/{repo_owner}/{repo_name}/archive/{archive_name}.zip") z = zipfile.ZipFile(io.BytesIO(r.content)) z.extractall(d) root_path = pathlib.Path(os.path.join(d, f"{repo_name}-{archive_name}")) docs_path = root_path.joinpath("docs/content", category) markdown_files = list(docs_path.glob("*.md*")) + list(docs_path.glob("*/*.md*")) for markdown_file in markdown_files: with open(markdown_file, "r") as f: relative_path = markdown_file.relative_to(root_path) github_url = f"https://github.com/{repo_owner}/{repo_name}/blob/{archive_name}/{relative_path}" yield Document(page_content=f.read(), metadata={"source": github_url})

然后,我们可以更新资产以获取源文档: @asset(compute_kind="GitHub") def source_docs(context: AssetExecutionContext): docs = [] for category in ["guides", "integrations"]: docs += list(get_github_docs("dagster-io", "dagster", category)) return docs

现在,让我们再次尝试“如何将 dbt 与 Dagster 配合使用?”:

好的,这个回答比之前相关性强多了。

优化性能与成本

然而,随着我们扩展管道的来源,成本可能会成为问题。

好消息是,管道会自动记录所有使用数据。让我们快速看一下:

好吧,成本略有增长。让我们深入挖掘!

看到成本增加促使我们进行优化。我们确信某些文档段落的更新频率低于其他部分,因此我们只需针对不同的文档部分在必要时更新嵌入和索引。

Dagster 的原生分区功能使我们能够有效地按类别组织源文档。以下是调整设置的方法:

我们首先定义一个 StaticPartitionsDefinition。让我们将每个文档类别分配为一个独立的分区: from dagster import ( StaticPartitionsDefinition, )

docs_partitions_def = StaticPartitionsDefinition(["guides", "integrations"])

然后,我们对资产进行分区: @asset(compute_kind="GitHub", partitions_def=docs_partitions_def) def source_docs(context: AssetExecutionContext): return list(get_github_docs("dagster-io", "dagster", context.partition_key))

@asset(compute_kind="OpenAI", partitions_def=docs_partitions_def) def search_index(context: AssetExecutionContext, openai: OpenAIResource, source_docs): source_chunks = [] splitter = CharacterTextSplitter(separator=" ", chunk_size=1024, chunk_overlap=0) for source in source_docs: context.log.info(source) for chunk in splitter.split_text(source.page_content): source_chunks.append(Document(page_content=chunk, metadata=source.metadata))

with openai.get_client(context) as client: search_index = FAISS.from_documents( source_chunks, OpenAIEmbeddings(client=client.embeddings) )

return search_index.serialize_to_bytes()

虽然我们的 completion 资产没有分区,但它依赖于新分区的 search_index 资产。使用分区映射可以允许 completion 资产依赖于 search_index 的所有分区: @asset( compute_kind="OpenAI", ins={ "search_index": AssetIn(partition_mapping=AllPartitionMapping()), }, ) def completion( context: AssetExecutionContext, openai: OpenAIResource, search_index: Dict[str, Any] ): merged_index: Any = None for index in search_index.values(): curr = FAISS.deserialize_from_bytes(index, OpenAIEmbeddings()) if not merged_index: merged_index = curr else: merged_index.merge_from(FAISS.deserialize_from_bytes(index, OpenAIEmbeddings())) question = "What can I use Dagster for?" with openai.get_client(context) as client: chain = load_qa_with_sources_chain(OpenAI(client=client.completions, temperature=0)) context.log.info( chain( { "input_documents": merged_index.similarity_search(question, k=4), "question": question, }, return_only_outputs=True, )["output_text"] )

现在,我们可以在资产图中看到 source_docs 和 search_index 资产各有 2 个分区:

现在,你可以刷新所需的分区,且仅在需要时刷新!

通过使用分区,我们可以在实际变更发生时,有选择地更新嵌入和搜索索引。这种方法最大限度地减少了不必要的计算,并有效降低了成本。

利用 LangChain 结合不同的 OpenAI 模型

好的,现在让我们探索不同的模型,以优化响应并获得更高的投资回报率。

我们可以参数化管道,以便根据查询选择最佳的 OpenAI 模型。方法如下: 首先,我们可以定义配置,以指定 OpenAI 模型及其应解决的问题: class OpenAIConfig(Config): model: str question: str

然后,让我们将此配置整合到 completion 资产中。 @asset( compute_kind="OpenAI", ins={ "search_index": AssetIn(partition_mapping=AllPartitionMapping()), }, )

def completion( context: AssetExecutionContext, openai: OpenAIResource, config: OpenAIConfig, search_index: Dict[str, Any], ): ... model = ChatOpenAI(client=client.chat.completions, model=config.model, temperature=0)

注意:为了动态应用不同的模型,我们使用 LangChain 的表达式语言(LCEL)来促进链的声明式组合。 这是一种声明式方法,用于真正组合链——并开箱即用地支持流式、批处理和异步处理。你可以使用所有现有的 LangChain 构建块来创建它们。 因此,我们的完整代码如下: @asset( compute_kind="OpenAI", ins={ "search_index": AssetIn(partition_mapping=AllPartitionMapping()), }, ) def completion( context: AssetExecutionContext, openai: OpenAIResource, config: OpenAIConfig, search_index: Dict[str, Any], ): merged_index: Any = None for index in search_index.values(): curr = FAISS.deserialize_from_bytes(index, OpenAIEmbeddings()) if not merged_index: merged_index = curr else: merged_index.merge_from(FAISS.deserialize_from_bytes(index, OpenAIEmbeddings())) with openai.get_client(context) as client: prompt = stuff_prompt.PROMPT model = ChatOpenAI(client=client.chat.completions, model=config.model, temperature=0) summaries = " ".join( [ SUMMARY_TEMPLATE.format(content=doc.page_content, source=doc.metadata["source"]) for doc in merged_index.similarity_search(config.question, k=4) ] ) context.log.info(summaries) output_parser = StrOutputParser() chain = prompt | model | output_parser context.log.info(chain.invoke({"summaries": summaries, "question": config.question}))

然后,我们可以手动使用 launchpad 启动运行:

对于问题“如何将 dbt 与 Dagster 配合使用?”,我们得到此答案:

到目前为止,我们的管道是手动启动的。

作为生产级管道,我们希望自动对新查询做出响应,无需人工干预。

为了演示目的,我们将传入的问题存储为指定文件夹内的 JSON 格式。 { "model": "gpt-3.5-turbo", "question": "How can I use dbt with Dagster?" }

以下是如何在 Dagster 中设置传感器,以便在新问题进入时自动物化答案: question_job = define_asset_job( name="question_job", selection=AssetSelection.keys(["completion"]), )

@sensor(job=question_job) def question_sensor(context): PATH_TO_QUESTIONS = os.path.join(os.path.dirname(__file__), "../../", "data/questions")

previous_state = json.loads(context.cursor) if context.cursor else {} current_state = {} runs_to_request = []

for filename in os.listdir(PATH_TO_QUESTIONS): file_path = os.path.join(PATH_TO_QUESTIONS, filename) if filename.endswith(".json") and os.path.isfile(file_path): last_modified = os.path.getmtime(file_path)

current_state[filename] = last_modified

if filename not in previous_state or previous_state[filename] != last_modified: with open(file_path, "r") as f: request_config = json.load(f)

runs_to_request.append( RunRequest( run_key=f"adhoc_request_{filename}_{last_modified}", run_config={"ops": {"completion": {"config": {request_config}}}}, ) )

return SensorResult(run_requests=runs_to_request, cursor=json.dumps(current_state))

注意:此设置主要为便于演示而设计。在实际场景中,将查询聚合到队列或数据库并批量处理到单次运行中,比为每个问题启动一个运行更高效。

传感器检测到我们的问题文件并启动运行以物化我们的 completion 资产。我们得到答案:

未来方向:提升 AI 管道的效率与效能

在这篇文章中,我们为优化 AI 奠定了基础,强调成本控制和提升开发者生产力。

关键要点包括:

- 推出新的 dagster-openai 集成,实现与 OpenAI API 的无缝交互以及开箱即用的使用情况追踪。 - 利用 LangChain 动态声明不同的模型。 - 利用现代编排器(Dagster)的功能,如 Insights、元数据追踪,以提高开发者生产力。

展望未来,我们计划进一步利用 Dagster 和 OpenAI 的能力进行模型性能比较和整体生产力提升。包括但不限于:

- 扩展使用 Dagster Cloud 的 Insights 功能,更全面地监控关键指标。 - 利用 Dagster 的数据目录,随着我们纳入更多来源(如集成 Slack 历史记录以改进响应质量)进行更好的管理。 - 确保数据新鲜度和可靠性,以维护支持机器人基础设施的质量,这涉及处理不良或伪造的来源。

有反馈或问题吗?请在 Slack 或 GitHub 上发起讨论。 有兴趣加入我们吗?查看我们的开放职位。 想要更多内容?在 LinkedIn 上关注我们。

评论 (0)