在实际企业级数据平台和 AI 应用开发中,Databricks 的 Lakehouse 架构已经成为处理大规模数据、构建机器学习流水线的关键基础设施。近期其市场动向也反映出,以 AI 智能体(AI Agent)为代表的新一代应用范式,正在成为驱动数据智能平台增长的核心引擎。对于开发者而言,理解如何利用 Databricks 这样的平台来构建、部署和管理 AI 智能体,是将前沿技术趋势转化为实际生产力的关键一步。
本文将从工程实践角度出发,探讨在 Databricks 环境中构建 AI 智能体的完整路径。我们将不局限于概念讨论,而是深入到环境配置、核心代码、工作流编排以及生产部署的各个环节。无论你是数据工程师希望为现有数据管道注入智能决策能力,还是机器学习工程师试图将大语言模型(LLM)与业务数据安全结合,或是应用开发者想要创建能够自主执行复杂任务的智能应用,本文提供的思路和示例都将为你提供一个清晰的起点。
1. 理解 AI 智能体与 Databricks 的协同价值
在深入代码之前,必须厘清两个核心概念:AI 智能体是什么,以及为什么 Databricks 是构建它的理想平台。这决定了后续所有技术选型和架构设计的方向。
1.1 AI 智能体:超越简单提示的自主系统
AI 智能体不是一个简单的聊天机器人或代码补全工具。它是一个具备感知、决策、执行和反思能力的软件实体。其核心特征在于能够理解高层次目标,并自主规划、调用工具(如搜索、计算、API)、执行动作来完成目标,过程中可以处理不确定性并从中学习。
一个典型的 AI 智能体工作流可能包括:
- 任务解析:将用户自然语言指令(如“分析上季度销售数据,找出下滑最严重的三个区域并给出原因”)分解为可执行的步骤。
- 工具调用:根据步骤,决定调用哪个工具。例如,调用 SQL 引擎查询销售数据,调用 Python 进行统计分析,调用外部 API 获取市场报告。
- 动作执行:在安全沙箱或授权环境下运行这些工具。
- 结果整合与输出:将各步骤的结果综合,生成结构化的报告、建议或执行下一步操作。
这与传统的数据分析或机器学习流水线有本质区别。传统流水线是预定义、确定性的,而智能体是目标驱动、动态规划的。
1.2 Databricks 作为智能体“大脑”的基座
Databricks 的 Lakehouse 平台融合了数据湖的灵活性和数据仓库的管理性,为 AI 智能体提供了独一无二的基础设施:
- 统一的数据治理:智能体需要访问大量企业数据。Databricks Unity Catalog 提供了统一的元数据、权限和审计层,确保智能体在合规的前提下访问数据,避免数据孤岛和权限混乱。
- 强大的计算引擎:无论是使用 Spark 进行 PB 级数据处理,还是使用单节点机器进行模型推理,Databricks 的弹性计算集群都能提供支撑。这对于需要处理海量上下文或进行复杂计算的智能体至关重要。
- 无缝的 ML 生命周期管理:MLflow 集成使得智能体所依赖的机器学习模型(包括 LLM)的版本管理、部署、监控变得标准化。智能体可以轻松调用生产环境中已部署的最佳模型。
- 安全与隔离:Databricks 工作区、集群和作业运行在受控的云环境中,为智能体的执行提供了网络隔离、身份认证和资源限制,降低了自主系统可能带来的风险。
将 AI 智能体构建在 Databricks 上,相当于为它配备了一个拥有海量记忆(数据)、强大算力(引擎)和严格行为准则(治理)的“数字大脑”。
2. 环境准备与核心依赖配置
构建一个生产可用的 AI 智能体,远不止安装一个 Python 库那么简单。我们需要在 Databricks 工作区中搭建一个兼顾开发效率和生产稳定的环境。
2.1 创建并配置 Databricks 集群
首先,你需要在 Databricks 工作区创建一个用于开发的集群。
集群配置选择:
- Databricks Runtime 版本:选择包含 ML 库的版本,如
13.3 LTS ML或更高。这预装了 TensorFlow、PyTorch 等常用库。 - 工作模式:对于开发和测试,
单用户模式即可。对于需要团队共享或作业调度的生产环境,考虑共享模式。 - 节点类型:根据智能体的复杂度选择。如果涉及大模型微调或复杂计算,选择 GPU 机型(如
g4dn.xlarge)。对于主要进行数据查询和轻量推理的智能体,CPU 机型(如i3.xlarge)更经济。 - 初始化脚本:这是关键一步。我们需要通过初始化脚本安装智能体框架所需的、Runtime 未预装的依赖。
- Databricks Runtime 版本:选择包含 ML 库的版本,如
编写集群初始化脚本: 在 Databricks 工作区的“计算”->“你的集群”->“配置”->“初始化脚本”中,添加一个脚本,例如
init_agent.sh。#!/bin/bash # 升级 pip 并设置镜像源(可选,用于加速国内下载) /databricks/python3/bin/pip install --upgrade pip /databricks/python3/bin/pip config set global.index-url https://pypi.tuna.tsinghua.edu.cn/simple # 安装核心智能体框架及工具依赖 # 这里以 LangChain 和 LlamaIndex 为例,它们是构建智能体的流行框架 /databricks/python3/bin/pip install langchain==0.1.0 langchain-community==0.0.10 /databricks/python3/bin/pip install llama-index==0.10.0 /databricks/python3/bin/pip install openai==1.3.0 # 如需调用 OpenAI API /databricks/python3/bin/pip install databricks-vectorsearch # Databricks 向量搜索(预览版) # 安装其他工具库,如 requests, sqlalchemy 等 /databricks/python3/bin/pip install requests sqlalchemy # 验证安装(可选) echo "初始化脚本执行完成,已安装智能体框架。"保存脚本后,需要将其路径(如
dbfs:/FileStore/scripts/init_agent.sh)配置到集群的初始化脚本设置中。重启集群以使安装生效。
2.2 配置密钥与权限
智能体通常需要调用外部 API(如 OpenAI)或访问数据库。永远不要将密钥硬编码在代码中。
使用 Databricks 密钥管理: 在 Databricks 工作区,进入“设置”->“密钥管理”。创建一个作用域(Scope),例如
agent-secrets。 在该作用域下添加密钥,例如openai-api-key,将你的 API 密钥值填入。在 Notebook 中安全读取密钥:
from databricks.sdk.runtime import dbutils # 安全地获取密钥 openai_api_key = dbutils.secrets.get(scope="agent-secrets", key="openai-api-key") # 现在可以在代码中使用 openai_api_key,它不会以明文形式出现在日志或代码中配置数据访问权限: 确保运行智能体的集群服务主体或用户,在 Unity Catalog 中拥有对目标数据表(如
sales_db.quarterly_report)的SELECT权限。这可以通过工作区的“数据”界面进行授权。
3. 构建一个与数据交互的初级智能体
现在,我们构建一个最简单的智能体:它接受一个关于销售数据的自然语言问题,将其转换为 SQL,在 Databricks 中执行,并返回答案。
3.1 项目结构与代码组织
在 Databricks Repos 中关联你的 Git 仓库,创建一个结构清晰的项目:
/my_ai_agent_project ├── README.md ├── requirements.txt # 生产依赖清单,内容与初始化脚本类似 ├── agent_core/ │ ├── __init__.py │ ├── sql_agent.py # SQL 智能体核心类 │ └── tools.py # 自定义工具定义 ├── notebooks/ │ └── 01_sql_agent_demo.ipynb # 交互式演示 Notebook └── tests/ └── test_agent.py3.2 实现 SQL 生成与执行工具
首先,在agent_core/tools.py中定义一个关键工具:执行 SQL 查询。
# agent_core/tools.py from langchain.tools import BaseTool from langchain.callbacks.manager import CallbackManagerForToolRun from typing import Optional, Type from pydantic import BaseModel, Field import pandas as pd from databricks import sql import os class SQLQueryInput(BaseModel): """SQL 查询工具的输入模型。""" query: str = Field(description="一个在 Databricks SQL 仓库中有效的、语法正确的 SQL 查询语句。") class DatabricksSQLTool(BaseTool): name = "databricks_sql_query" description = "在指定的 Databricks SQL 仓库上执行一个 SQL 查询,并返回结果。用于回答关于销售、用户或产品数据的问题。输入必须是一个完整的 SQL 语句。" args_schema: Type[BaseModel] = SQLQueryInput return_direct: bool = False # 让 Agent 决定下一步 def _run(self, query: str, run_manager: Optional[CallbackManagerForToolRun] = None) -> str: """执行查询并返回格式化结果。""" # 从 Databricks 上下文中获取连接参数(更安全的方式) # 假设服务器主机名、HTTP路径、令牌已通过集群配置或环境变量设置 server_hostname = spark.conf.get("spark.databricks.workspaceUrl") http_path = spark.conf.get("spark.databricks.sql.warehouse.httpPath") access_token = dbutils.secrets.get(scope="agent-secrets", key="databricks-token") if not all([server_hostname, http_path, access_token]): return "错误:未正确配置 Databricks SQL 仓库连接信息。请检查集群配置或环境变量。" connection = sql.connect( server_hostname=server_hostname, http_path=http_path, access_token=access_token ) cursor = connection.cursor() try: cursor.execute(query) result = cursor.fetchall() # 将结果转换为 Pandas DataFrame 以便于格式化 columns = [desc[0] for desc in cursor.description] df = pd.DataFrame(result, columns=columns) # 返回前N行作为字符串,避免数据过大 return f"查询成功。结果(前10行):\n{df.head(10).to_string(index=False)}" except Exception as e: return f"执行 SQL 查询时出错:{str(e)}" finally: cursor.close() connection.close()3.3 组装智能体并定义系统提示词
接下来,在agent_core/sql_agent.py中创建智能体。
# agent_core/sql_agent.py from langchain.agents import create_react_agent, AgentExecutor from langchain.tools import Tool from langchain_openai import ChatOpenAI from langchain.prompts import PromptTemplate from .tools import DatabricksSQLTool import os def create_sql_agent(openai_api_key: str): """ 创建一个专用于 SQL 查询的智能体。 """ # 1. 初始化 LLM llm = ChatOpenAI( model="gpt-4", # 或 "gpt-3.5-turbo",GPT-4 在复杂逻辑和 SQL 生成上更可靠 temperature=0, # 对于 SQL 生成,低温度确保确定性 openai_api_key=openai_api_key ) # 2. 准备工具列表 sql_tool = DatabricksSQLTool() tools = [sql_tool] # 3. 定义系统提示词 - 这是智能体行为的核心 prompt_template = """ 你是一个专业的 Databricks SQL 分析师。你的任务是根据用户的问题,生成正确的 SQL 查询语句,并使用 `databricks_sql_query` 工具执行它,最后将结果以清晰、易懂的语言总结给用户。 请遵循以下步骤思考(Thought/Action/Observation): 1. 理解用户的问题,确定需要查询的数据表和字段。 2. 根据已知的数据库 schema(例如,我们有 `sales.orders` 表,包含 `order_id`, `region`, `amount`, `order_date` 等字段),编写一个语法正确的 Databricks SQL 查询。 3. 只使用 `databricks_sql_query` 这一个工具。你的 Action 必须是调用这个工具,并将完整的 SQL 语句作为输入。 4. 观察工具返回的结果。如果结果是错误信息,分析错误并尝试修正 SQL 语句。 5. 将最终的结果用自然语言总结给用户,可以包含关键数字和趋势。 已知数据库 Schema 摘要: - 数据库: `sales` - 表: `orders` (订单表) - `order_id` (整数,主键) - `customer_id` (整数) - `region` (字符串,如 'North', 'South') - `amount` (小数,订单金额) - `order_date` (日期) - 表: `products` (产品表) - `product_id` (整数) - `product_name` (字符串) - `category` (字符串) 现在开始。如果用户的问题无法通过查询现有数据回答,请礼貌说明。 问题:{input} 思考过程: """ prompt = PromptTemplate.from_template(prompt_template) # 4. 使用 ReAct 框架创建智能体 agent = create_react_agent(llm=llm, tools=tools, prompt=prompt) # 5. 创建执行器 agent_executor = AgentExecutor( agent=agent, tools=tools, verbose=True, # 开发时设为 True 以查看思考链 handle_parsing_errors=True, # 优雅处理解析错误 max_iterations=5, # 防止无限循环 early_stopping_method="generate" ) return agent_executor3.4 在 Notebook 中运行与测试
在notebooks/01_sql_agent_demo.ipynb中,我们可以进行交互式测试。
# 单元格 1:设置环境并导入 import sys sys.path.append('/Workspace/Repos/your_repo/my_ai_agent_project') # 添加项目路径 from agent_core.sql_agent import create_sql_agent from databricks.sdk.runtime import dbutils # 安全获取密钥 openai_api_key = dbutils.secrets.get(scope="agent-secrets", key="openai-api-key") # 创建智能体 agent = create_sql_agent(openai_api_key) # 单元格 2:运行测试 question = "上个季度,哪个区域的销售总额最高?是多少?" result = agent.invoke({"input": question}) print(result["output"])当verbose=True时,你将在输出中看到类似以下的思考链,这对于调试和理解智能体行为至关重要:
> 进入新的 AgentExecutor 链... 思考:用户想知道上个季度哪个区域销售总额最高。我需要查询 `sales.orders` 表,按区域汇总金额,并过滤出上个季度的数据。首先,我需要确定上个季度的日期范围。 行动:{ "action": "databricks_sql_query", "action_input": "SELECT region, SUM(amount) as total_sales FROM sales.orders WHERE order_date >= '2024-01-01' AND order_date < '2024-04-01' GROUP BY region ORDER BY total_sales DESC LIMIT 1" } 观察:查询成功。结果(前10行): region total_sales North 1250000.50 思考:工具返回了结果。最高销售额的区域是 North,总额为 1,250,000.50。我需要将这个信息总结给用户。 最终答案:根据查询结果,上个季度(2024年第一季度)销售额最高的区域是 **North(北部)**,销售总额为 **1,250,000.50**。 > 链结束。4. 进阶:构建具备多工具与记忆能力的智能体
初级 SQL 智能体功能单一。一个真正的智能体需要组合多种工具,并具备记忆能力以进行多轮对话。
4.1 扩展工具集
在tools.py中增加更多工具,例如:
- 数据可视化工具:调用
matplotlib或plotly生成图表,将图片保存到 DBFS 并返回路径。 - 文件读取工具:从 DBFS 或云存储中读取 CSV、JSON 文件。
- API 调用工具:获取天气、汇率等外部信息。
- Python 代码执行工具:在安全沙箱中执行计算或数据处理(需极其谨慎,避免任意代码执行风险)。
每个工具都需要明确定义name,description和args_schema。清晰的描述能帮助 LLM 准确选择工具。
4.2 为智能体添加记忆
LangChain 提供了多种记忆后端。对于对话式智能体,ConversationBufferWindowMemory是一个简单选择。
# agent_core/chat_agent.py from langchain.memory import ConversationBufferWindowMemory from langchain.agents import AgentExecutor, create_react_agent from langchain_openai import ChatOpenAI from .tools import DatabricksSQLTool, DataVizTool # 假设我们新增了可视化工具 def create_chat_agent_with_memory(openai_api_key: str): llm = ChatOpenAI(model="gpt-4", temperature=0.1, openai_api_key=openai_api_key) tools = [DatabricksSQLTool(), DataVizTool()] # 创建记忆,保留最近3轮对话 memory = ConversationBufferWindowMemory( memory_key="chat_history", k=3, return_messages=True ) # 提示词模板需要包含 `chat_history` 和 `input` 两个变量 prompt_template = """ 你是一个数据分析助手,可以查询数据、生成图表。你有以下工具: {tools} 使用以下格式: 历史对话: {chat_history} 问题:{input} 思考:我需要先分析用户的问题... 行动:{agent_scratchpad} """ prompt = PromptTemplate.from_template(prompt_template) agent = create_react_agent(llm=llm, tools=tools, prompt=prompt) agent_executor = AgentExecutor.from_agent_and_tools( agent=agent, tools=tools, memory=memory, verbose=True, handle_parsing_errors=True ) return agent_executor这样,智能体就能在对话中引用之前的上下文,例如用户问“那么它的趋势如何?”,智能体能知道“它”指的是上一轮查询的数据。
5. 生产部署与运维考量
在 Notebook 中交互运行只是第一步。要让智能体服务化,并稳定运行于生产环境,需要考虑以下方面。
5.1 部署为 Databricks 作业或模型服务
作业部署:将智能体逻辑封装在一个 Python 脚本中,通过 Databricks 作业按计划或触发式运行。适合处理批量任务(如每日自动生成报告)。
- 创建作业,配置集群、依赖库。
- 将主脚本路径指向你的
main.py。 - 在作业参数中传递运行指令(如
--task “analyze_sales_trend”)。
模型服务部署:如果智能体需要提供低延迟的 API 服务(如聊天接口),可以使用 MLflow 和 Databricks Model Serving。
- 使用 MLflow 的
pyfunc模型格式包装你的智能体类。 - 记录(log)模型到 MLflow Registry。
- 将注册的模型部署为实时服务端点(Serverless 或 Classic)。
- 通过 REST API 调用智能体。
- 使用 MLflow 的
5.2 监控、日志与错误处理
- 结构化日志:使用 Python
logging模块,记录智能体的关键决策点、工具调用、LLM 请求和最终输出。将日志发送到集中式日志系统(如 Datadog, Splunk)。 - 性能监控:监控每个请求的端到端延迟、Token 消耗(成本)、工具调用成功率。
- 错误隔离与重试:为工具调用(如 SQL 查询、API 调用)添加重试逻辑和断路器模式。确保单个工具失败不会导致整个智能体崩溃,并能向用户返回友好的错误信息。
- LLM 输出验证:对于关键操作(如执行数据删除的 SQL),在智能体执行前,可以设计一个“确认”步骤,或通过另一套规则引擎验证 LLM 生成的指令是否安全。
5.3 成本与性能优化
- LLM 调用优化:
- 使用缓存(如
LangChain的InMemoryCache或RedisCache)存储重复问题的结果。 - 对长上下文进行摘要或选择性检索,减少提示词 Token 数。
- 根据任务复杂度选择模型,简单任务用
gpt-3.5-turbo,复杂任务再用gpt-4。
- 使用缓存(如
- 向量检索集成:对于需要基于大量文档知识回答的问题,可以将公司文档存入 Databricks Vector Search,让智能体先检索相关片段,再基于片段生成答案,提高准确性和降低幻觉。
6. 常见问题排查清单
在开发和运行 AI 智能体时,你会遇到各种问题。下表列出了常见现象、原因及排查方向。
| 问题现象 | 可能原因 | 排查步骤 | 解决方案 |
|---|---|---|---|
| 智能体无法选择正确的工具 | 1. 工具描述 (description) 不清晰或与用户问题不匹配。2. LLM 温度 ( temperature) 设置过高,导致输出随机。3. 系统提示词未明确指导工具使用。 | 1. 检查verbose=True时的思考链,看 LLM 是如何理解问题和工具的。2. 简化工具描述,使其更精准。 3. 在提示词中举例说明工具使用场景。 | 1. 重写工具描述,使其功能一目了然。 2. 将 temperature设为 0 或接近 0。3. 在提示词中加入少量示例(Few-shot)。 |
| SQL 查询执行报语法错误 | 1. LLM 生成的 SQL 不符合 Databricks 方言。 2. 表名或列名不正确。 3. 日期等字面量格式错误。 | 1. 在系统提示词中明确说明使用“Databricks SQL”语法。 2. 提供更准确的 Schema 信息。 3. 捕获工具错误并打印完整 SQL。 | 1. 在提示词中提供更详细的 Schema,甚至包含示例查询。 2. 在工具层增加一个 SQL 语法预检(使用轻量级解析库)。 3. 让智能体在错误后尝试修正(ReAct 框架支持)。 |
| 智能体陷入循环或超时 | 1.max_iterations设置过高。2. 工具执行结果未能让智能体达到“最终答案”状态。 3. 提示词未定义明确的停止条件。 | 1. 查看verbose日志,观察智能体在重复什么动作。2. 检查工具返回的结果格式是否清晰。 | 1. 合理设置max_iterations(如 5-10)。2. 优化工具返回结果,使其包含明确的任务完成信号。 3. 在提示词中强调“当你得到最终答案时,就直接回复用户,停止使用工具”。 |
| 调用外部 API 失败 | 1. 网络不通(集群无外网权限或 VPC 限制)。 2. API 密钥未正确配置或已失效。 3. 请求频率超限。 | 1. 在 Notebook 中手动运行requests.get(api_url)测试连通性。2. 检查密钥是否通过 dbutils.secrets.get正确获取。3. 查看 API 提供商的监控面板。 | 1. 配置集群的 VPC、安全组和 NAT 网关以允许出站流量。 2. 确保密钥作用域和名称正确,并已授予 Notebook/作业访问权限。 3. 在代码中添加退避重试机制。 |
| 内存消耗过大或进程崩溃 | 1. 处理的数据量过大(如 SQL 返回百万行)。 2. LLM 上下文过长。 3. 智能体迭代次数太多,累积中间结果。 | 1. 监控集群的 Spark UI 或 Ganglia 指标。 2. 检查工具返回的数据大小。 | 1. 在 SQL 工具中强制增加LIMIT子句或进行预聚合。2. 对长文本进行分块或摘要后再喂给 LLM。 3. 定期清理记忆或使用摘要式记忆。 |
7. 最佳实践与扩展方向
基于项目经验,遵循以下实践能显著提升智能体的可靠性和价值。
7.1 安全与治理第一
- 最小权限原则:为智能体使用的服务主体分配最小必要的数据和 API 访问权限。使用 Unity Catalog 进行细粒度数据管控。
- 输入输出净化:对用户输入进行基本的清理和检查,防止提示词注入攻击。对智能体生成的、将要被执行的代码或命令进行严格的沙箱测试或人工审核(特别是对于写操作)。
- 审计日志:记录所有用户查询、智能体思考过程、工具调用详情和最终输出,便于事后追溯和分析。
7.2 设计可评估的智能体
- 定义明确的成功指标:不仅仅是“能运行”,而是“准确率”(生成的 SQL 正确率)、“完成率”(成功解答问题的比例)、“用户满意度”等。
- 构建测试集:针对不同业务场景,构建一批标准问题,定期运行智能体进行回归测试,监控其性能变化。
- A/B 测试:对比不同提示词、不同模型(如 GPT-4 vs. Claude)或不同工作流设计的效果。
7.3 走向更复杂的智能体系统
当前示例是一个单体智能体。更复杂的生产系统可能需要:
- 多智能体协作:设计专精于不同领域(查询、分析、报告)的智能体,让它们通过消息队列或协调器协同完成复杂任务。
- 与向量数据库深度集成:利用 Databricks Vector Search,让智能体具备检索企业内部知识库(产品文档、会议纪要)的能力,实现基于知识的问答(RAG)。
- 长期记忆与个性化:将对话历史、用户偏好存储到数据库中,使智能体能够提供个性化的、有连续上下文的服务。
构建 AI 智能体是一个迭代过程,从解决一个明确、具体的小问题开始(如自动生成周报 SQL),逐步扩展其能力和边界。Databricks 提供的统一数据、治理和计算平台,极大地简化了将数据与智能体结合的基础设施复杂度,让开发者能更专注于智能体本身的逻辑与创新。