本指南提供 Drasi 工具的快速入门概述。如需了解所有 Drasi 功能、参数和配置的详细列表,请访问 Drasi 文档,以及 langchain_drasi repository.
概述
Drasi 是一个变更检测平台,可以轻松高效地检测和响应数据库中的变更。LangChain-Drasi 集成通过将外部数据变更与工作流执行相连接,创建响应式、变更驱动的 AI 代理。这使得代理能够通过将外部数据变更与代理工作流相桥接,从而发现、订阅和响应实时查询更新。Drasi 连续查询流式传输实时更新,触发代理状态转换、修改内存或动态控制工作流执行——将静态代理转变为环境持久、响应迅速的系统。
详细信息
| 类 | 包 | 可序列化 | JS 支持 | 下载量 | 版本 |
|---|---|---|---|---|---|
DrasiTool | langchain-drasi | ❌ | ❌ | !PyPI - 下载量 | !PyPI - 版本 |
功能特性
- 查询发现 - 自动识别可用的 Drasi 查询
- 实时订阅 - 监控连续查询更新
- 通知处理器 - 六种内置处理器,适用于不同场景
- - 控制台
- - 日志记录
- - 内存
- - 缓冲区
- - LangChain 内存
- - LangGraph 内存
- 自定义处理器 - 扩展基础处理器以实现领域特定逻辑
设置
要访问 Drasi 工具,您需要运行 Drasi 和 Drasi MCP 服务器。
前置条件
- - Drasi 平台 - 已安装并运行
- - Drasi MCP 服务器 - 已配置并可访问
- - Python 3.11+ - 需要
langchain-drasi包
凭证(可选)
如果您的 Drasi MCP 服务器需要身份验证,您可以使用 Bearer 令牌或其他身份验证方法配置标头:
from langchain_drasi import MCPConnectionConfig
config = MCPConnectionConfig(
server_url="http://localhost:8083",
headers={"Authorization": "Bearer your-token"},
timeout=30.0
)
安装
Drasi 工具位于 langchain-drasi package:
pip install -U langchain-drasi
uv add langchain-drasi
实例化
现在我们可以实例化 Drasi 工具。您需要配置 MCP 连接,并可选择添加通知处理器来处理实时更新:
from langchain_drasi import create_drasi_tool, MCPConnectionConfig, ConsoleHandler
# Configure connection to Drasi MCP server
config = MCPConnectionConfig(
server_url="http://localhost:8083",
timeout=30.0
)
# Create a notification handler
handler = ConsoleHandler()
# Create the tool
tool = create_drasi_tool(
mcp_config=config,
notification_handlers=[handler]
)
调用
直接调用
以下是直接调用该工具的简单示例。
# Discover available queries
queries = await tool.discover_queries()
# Returns: [QueryInfo, QueryInfo, ...]
# Subscribe to a specific query
await tool.subscribe("hot-freezers")
# Notifications routed to registered handlers
# Read current results from a query
result = await tool.read_query("active-orders")
# Returns: QueryResult with current data
As a ToolCall
我们还可以使用模型生成的 ToolCall,在这种情况下ToolMessage将返回。
在智能体中
我们可以在 LangGraph 智能体中使用 Drasi 工具来创建响应式的事件驱动工作流。为此,我们需要一个具有工具调用能力的模型。
from langchain_anthropic import ChatAnthropic
from langchain.agents import create_agent
# Initialize the model
model = ChatAnthropic(model="claude-sonnet-4-6")
# Create agent with Drasi tool
agent = create_agent(model, [tool])
# Run the agent
result = agent.invoke(
{"messages": [{"role": "user", "content": "What queries are available?"}]}
)
print(result["messages"][-1].content)
result = agent.invoke(
{"messages": [{"role": "user", "content": "Subscribe to the customer-orders query"}]}
)
print(result["messages"][-1].content)
通知处理器
Drasi 的关键特性之一是其内置的通知处理器,用于处理实时查询结果变化。您可以使用这些处理器根据数据变化执行特定操作。
内置处理器
ConsoleHandler - 将格式化的通知输出到标准输出:
from langchain_drasi import ConsoleHandler
handler = ConsoleHandler()
LoggingHandler - 使用 Python 的日志框架记录通知:
from langchain_drasi import LoggingHandler
handler = LoggingHandler(
logger_name="drasi.notifications",
log_level=logging.INFO
)
MemoryHandler - 在内存中存储通知,支持可选的过滤功能:
from langchain_drasi import MemoryHandler
handler = MemoryHandler(max_size=100)
# Retrieve notifications
all_notifs = handler.get_all()
freezer_notifs = handler.get_by_query("hot-freezers")
added_events = handler.get_by_type("added")
BufferHandler - 用于顺序处理的先进先出队列:
这对于在工作流忙于处理其他任务时缓冲传入的变更通知非常有用;您可以在工作流中设置一个循环,以便在就绪时从缓冲区消费通知。
from langchain_drasi import BufferHandler
handler = BufferHandler(max_size=100)
# Later, consume notifications
notification = handler.consume() # Remove and return next notification
notification = handler.peek() # View next notification without removing
LangGraphMemoryHandler - 将更新直接注入 LangGraph 检查点:
from langchain_drasi import LangGraphMemoryHandler
from langgraph.checkpoint.memory import MemorySaver
checkpoint_manager = MemorySaver()
handler = LangGraphMemoryHandler(
checkpointer=checkpoint_manager,
thread_id="your-thread-id"
)
自定义处理器
您可以通过扩展来创建自定义处理器 BaseDrasiNotificationHandler:
from langchain_drasi import BaseDrasiNotificationHandler
class CustomHandler(BaseDrasiNotificationHandler):
def on_result_added(self, query_name: str, added_data: dict):
# Handle new results
print(f"New result in {query_name}: {added_data}")
def on_result_updated(self, query_name: str, updated_data: dict):
# Handle updated results
print(f"Updated result in {query_name}: {updated_data}")
def on_result_deleted(self, query_name: str, deleted_data: dict):
# Handle deleted results
print(f"Deleted result in {query_name}: {deleted_data}")
handler = CustomHandler()
tool = create_drasi_tool(
mcp_config=config,
notification_handlers=[handler]
)
示例
使用场景
Drasi 特别适用于构建需要响应实时数据变化的环境智能体。一些示例使用场景包括:
- AI 副驾驶 - 监控并响应系统事件的助手
- AI 游戏玩家 - 适应游戏事件的 NPC
- 物联网监控 - 处理传感器数据流的智能体
- 客户支持 - 对工单更新或客户操作做出响应的机器人
- DevOps 助手 - 监控基础设施变化的工具
- 协作编辑 - 响应文档或代码变化的系统
API 参考
有关所有 Drasi 功能和配置的详细文档,请参阅 API 参考.