
Apache Airflow LangChainHook 实战指南用 Airflow Connection 统一管理 LangChain 对话与嵌入模型【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文以 Apache Airflow 的apache-airflow-providers-common-aiProvider 中的LangChainHook为主线讲解如何用一条 Airflow ConnectionAPI Key、可选 Base URL桥接 LangChain 的init_chat_model与init_embeddings两个通用入口统一提供对话模型与嵌入模型对象。读完本文你将掌握langchain连接类型的字段配置、模型标识解析优先级、单实例双模型与双连接隔离的写法以及从源码层面理解该 Hook 的凭证转发机制与连接测试逻辑。LangChainHook 的定位与工作方式LangChainHook是 Airflow 与 LangChain 生态之间的“翻译层”它把 Airflow 的 Connection存放 API Key 与可选 Base URL翻译成 LangChain 模型构造函数所需的api_key/base_url参数并通过两个厂商无关的通用入口完成厂商分发langchain.chat_models.init_chat_model负责对话模型根据provider:name前缀如openai:gpt-4o分发到对应厂商实现langchain.embeddings.init_embeddings负责嵌入模型同样的分发机制。与 LlamaIndex、Pydantic AI 等 Provider 内其他 Hook 不同LangChainHook拥有自己的langchain连接类型这样 Airflow UI 中的 Connections 表单能如实反映“这条连接配置的是哪个框架”。从 provider.yaml 可以看到该 Hook 的完整注册信息包括hook-class-name: airflow.providers.common.ai.hooks.langchain.LangChainHook、connection-type: langchain、外部服务清单OpenAI、Anthropic、Groq、Mistral AI、DeepSeek、Ollama、vLLM以及连接表单的 UI 行为定义。Hook 的核心实现位于 hooks/langchain.py其中几个类属性决定了它的身份conn_name_attr llm_conn_id default_conn_name langchain_default conn_type langchain hook_name LangChain对话模型Chat Model用法在构造函数中传入llm_model或在 Connection 的extra[model]中设置然后调用get_chat_model()即可拿到一个 LangChainBaseChatModel。仓库中的示例 DAG example_langchain_hook.py 给出了标准写法dag(scheduleNone, tags[example]) def example_langchain_chat(): task def summarize(text: str) - str: hook LangChainHook( llm_conn_idlangchain_default, llm_modelopenai:gpt-4o, ) llm hook.get_chat_model() # LangChain BaseMessage.content is str | list[...] (multi-modal union); # coerce to str for the text-only path this example demonstrates. return str(llm.invoke(fSummarize concisely: {text}).content) summarize(Apache Airflow is a platform for authoring, scheduling, and monitoring workflows.)由于返回值是标准的BaseChatModel它可以直接与 LangChain 的 Runnable 体系组合ChatPromptTemplate/StrOutputParser/RunnableSequence等无需任何额外适配。从源码看get_chat_model()的执行链非常清晰见 hooks/langchain.py惰性导入from langchain.chat_models import init_chat_model——langchain是可选 extra模块级导入会让未安装该 extra 的用户无法加载整个common.ai包通过self.get_connection(self.llm_conn_id)取出 Airflow Connection调用_resolve_model_id()从“构造参数或 Connection extra”中解析出模型标识执行init_chat_model(model_id, **self._connection_kwargs(conn))完成分发。其中_connection_kwargs的凭证映射规则是见 源码Connection 的password非空 → 转发为api_keyConnection 的host非空 → 转发为base_url两者都为空时不传任何额外参数例如本地无鉴权场景。支持的对话模型厂商任何被init_chat_model接受的模型标识都可以直接使用常见标识如下每个厂商需要对应安装 LangChain 集成包openai:gpt-4o、openai:gpt-4o-mini—— 需要langchain-openaianthropic:claude-sonnet-5—— 需要langchain-anthropicgroq:llama-3.3-70b-versatile—— 需要langchain-groqmistralai:mistral-large-latest—— 需要langchain-mistralaiollama:llama3—— 需要langchain-ollama把host指向 Ollama 的 URLdeepseek:deepseek-chat—— 需要langchain-deepseek需要注意的是采用非标准鉴权的云厂商AWS Bedrock、Google Vertex AI、Azure OpenAI不在api_keybase_url的凭证面覆盖范围内文档将其留给按厂商划分的专用 Hook与 pydantic-ai 的PydanticAIBedrockHook/PydanticAIVertexHook/PydanticAIAzureHook子类模式保持一致。从源码结构看这些云厂商 Hook 确实已在 hooks/pydantic_ai.py 中按子类方式实现LangChain 侧未来若扩展也将沿用同一模式。嵌入模型Embedding Model用法在构造函数中传入embed_model或设置extra[embed_model]调用get_embedding_model()即可获得Embeddings对象dag(scheduleNone, tags[example]) def example_langchain_embedding(): task def embed_documents(texts: list[str]) - int: hook LangChainHook( llm_conn_idlangchain_default, embed_modelopenai:text-embedding-3-small, ) embeddings hook.get_embedding_model() vectors embeddings.embed_documents(texts) return len(vectors[0]) embed_documents( [ Apache Airflow is a workflow orchestrator., Workflows are defined as Python DAGs., ] )单实例同时服务对话与嵌入当两类模型标识都配置好时同一个 Hook 实例可以复用避免为同一厂商重复构造连接dag(scheduleNone, tags[example]) def example_langchain_chat_and_embedding(): One hook instance serves both chat and embeddings when both models are set. task def use_both() - dict: hook LangChainHook( llm_conn_idlangchain_default, llm_modelopenai:gpt-4o, embed_modelopenai:text-embedding-3-small, ) chat hook.get_chat_model() embeddings hook.get_embedding_model() return { answer: str(chat.invoke(In one sentence: what does Airflow do?).content), embedding_dim: len(embeddings.embed_query(Airflow)), } use_both()支持的嵌入模型厂商Hook 会把连接中的api_key与可选base_url传给init_embeddings因此只有接受这种参数形态的嵌入类才能直接工作openai:text-embedding-3-small、openai:text-embedding-3-large—— 需要langchain-openai面向 OpenAI 兼容端点Ollama / vLLM / LM Studio的openai:model—— 同样需要langchain-openai把host指向对应端点即可。init_embeddings虽然还声明了更多厂商Cohere、Mistral AI、HuggingFace、Bedrock、Vertex AI、Azure OpenAI 等但它们的嵌入类期望的是厂商专属凭证参数cohere_api_key、AWS 鉴权链、GCP 服务账号等与 Hook 转发的通用api_key/base_url不匹配因此不在当前连接类型覆盖范围内同样留给按厂商划分的子类。Chat 与 Embeddings 使用不同连接如果对话与嵌入走不同的 API Key例如付费对话 Key 免费额度的嵌入 Key传入显式的embed_conn_id不设置时它自动回落到llm_conn_id单厂商常见场景保持最简dag(scheduleNone, tags[example]) def example_langchain_different_conns(): Use separate connections when chat and embeddings live on different API keys. task def use_separate_conns() - dict: hook LangChainHook( llm_conn_idopenai_chat, embed_conn_idopenai_embed, llm_modelopenai:gpt-4o, embed_modelopenai:text-embedding-3-small, ) chat hook.get_chat_model() embeddings hook.get_embedding_model() return { answer: str(chat.invoke(In one sentence: what does Airflow do?).content), embedding_dim: len(embeddings.embed_query(Airflow)), } use_separate_conns()这一点在单元测试 test_langchain.py 中有专门验证TestLangChainHookInit::test_embed_conn_falls_back_to_llm_conn断言未传embed_conn_id时回落为llm_conn_idTestGetEmbeddingModel::test_uses_embed_conn_id_when_set则验证设置后get_embedding_model()会读取embed_conn的独立 API Key 与模型。Connection 配置详解LangChainHook从连接类型为langchain的 Airflow Connection 读取凭证各字段含义如下password—— API Key转发为init_chat_model/init_embeddings的api_key参数host—— 可选的 Base URL转发为base_url参数适用于自定义 OpenAI 兼容端点、Ollama、vLLM 等自托管场景extra JSON——{model: openai:gpt-4o, embed_model: openai:text-embedding-3-small}在连接上设置默认的对话与嵌入模型标识。另外schema、port、login三个字段在连接表单中被隐藏get_ui_field_behaviour()返回hidden_fields: [schema, port, login]并把password重命名为 “API Key”本连接类型不使用它们。provider.yaml 中还为该连接定义了conn-fieldsmodel与embed_model使连接表单中出现专门的输入项其值最终仍存储在extra[model]/extra[embed_model]中。配套的连接文档见 connections/langchain.rst其中给出了几种典型连接 JSON{ conn_type: langchain, password: sk-..., extra: {\model\: \openai:gpt-4o\, \embed_model\: \openai:text-embedding-3-small\} }{ conn_type: langchain, password: sk-ant-..., extra: {\model\: \anthropic:claude-sonnet-5\} }{ conn_type: langchain, host: http://localhost:11434/v1, extra: {\model\: \ollama:llama3\} }模型标识解析优先级get_chat_model()与get_embedding_model()解析模型标识的顺序一致实现于_resolve_model_id()构造参数llm_model/embed_model优先级最高覆盖连接 extra连接上的extra[model]/extra[embed_model]。两处都未设置时调用模型方法将抛出ValueError错误信息会明确提示“Passmodel/embed_model to the hook constructor or set extra on the connection”。单元测试TestResolveModelId覆盖了三类断言构造参数优先于 extra、extra 兜底生效、两处皆空时分别抛出No chat model identifier set/No embedding model identifier set。test_connection只做解析不真正调用 LLMLangChainHook实现了test_connection()见 源码供 Airflow UI 的 “Test” 按钮使用。它的语义是校验模型标识有效、且厂商类能用给定凭证完成实例化成功返回(True, Model resolved successfully.)失败时返回异常信息。值得强调的是其设计取舍它不发起真实的 LLM API 调用——那样既昂贵又会因配额、计费、限流等与连通性无关的原因而误报失败。单元测试TestConnectionTest验证了三种情形成功解析、init_chat_model抛出Unknown provider时返回失败及错误文本、未设置模型时返回No chat model identifier set。参数一览参数默认值说明llm_conn_idlangchain_defaultLLM 厂商对应的 Airflow 连接 ID。embed_conn_idNone回落到llm_conn_id嵌入厂商可选的独立连接 ID。适用于对话与嵌入使用不同 API Key 的场景单厂商场景留空即可Hook 会复用llm_conn_id。llm_modelNone回落到连接的extra[model]provider:name形式的对话模型标识如openai:gpt-4o。仅在调用get_chat_model()时需要。embed_modelNone回落到连接的extra[embed_model]provider:name形式的嵌入模型标识如openai:text-embedding-3-small。仅在调用get_embedding_model()时需要。从源码细节看__init__中llm_conn_id的默认值是在运行时用self.default_conn_name解析的而非直接绑定类属性——源码注释解释了原因Python 在类定义时求值默认参数直接写llm_conn_id: str default_conn_name会把基类值固化导致未来定义自定义default_conn_name的子类无法继承正确默认值见 hooks/langchain.py。依赖安装使用本 Hook 需安装langchainextrapip install apache-airflow-providers-common-ai[langchain]查看 pyproject.toml 可知该 extra 仅声明langchain1.0.0——框架本身是厂商无关的因此不会替你安装任何厂商包。你还需要自行安装目标厂商的 LangChain 集成包langchain-openai—— OpenAI 及 OpenAI 兼容端点Ollama、vLLM 等langchain-anthropic—— Anthropiclangchain-groq、langchain-mistralai、langchain-deepseek、langchain-ollama等 —— 对应其余厂商这种“extra 只装框架、厂商包按需添加”的拆分与 Hook 内部的惰性导入设计相呼应未安装langchain的用户导入common.ai其他模块不受影响只有真正调用get_chat_model()/get_embedding_model()时才会触发langchain导入。小结LangChainHook以极小的接口面两个连接 ID、两个模型标识覆盖了 Airflow 任务中使用 LangChain 对话与嵌入模型的常见路径凭证统一收口到langchain类型的 Connection模型分发交给init_chat_model/init_embeddings自托管端点通过host→base_url一行映射接入 Ollama / vLLM。其适用边界也很明确——只服务api_keybase_url凭证面的厂商Bedrock、Vertex AI、Azure OpenAI 等需要专属鉴权参数的厂商不在本连接类型覆盖范围内需要时请参考 Provider 内 pydantic-ai 系列的按厂商子类实现hooks/pydantic_ai.py。完整可运行的示例 DAG 位于 example_dags/example_langchain_hook.py行为验证可参考 tests/unit/common/ai/hooks/test_langchain.py。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考