ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Doris + MCP:Agent时代实时数据分析的黄金组合(技术解析+实战案例)

Apache Doris + MCP:Agent时代实时数据分析的黄金组合(技术解析+实战案例) 1. 为什么 Agent 需要 Apache Doris MCP 这套组合Apache Doris 是一款开源的 MPP 分析型数据库主打实时写入与亚秒级查询MCPModel Context Protocol是一套让大模型与外部数据源、工具标准化对话的协议。把两者接在一起本质上是给 Agent 装上一个能实时查数、还能自己发现有哪些表可查的数据底座。它适合谁适合正在做 ChatBI、智能报表、风控巡检、IoT 监控看板又不想为每个数据源手写适配层的后端和算法同学。传统做法里Agent 想查数据库要么把 SQL 硬编码进 prompt要么给每个库写一套 Function Calling 描述。表结构一变工具描述就得跟着改维护成本高得离谱。更麻烦的是延迟批处理报表 T1 才出数Agent 拿到的是昨天的世界做实时决策根本不够用。我见过一个库存场景运营问现在哪个 SKU 快断货了Agent 只能回昨天的快照等报表刷新爆款早卖空了。Apache Doris 解决的是算得快、写得进这一层。它用 MPP 架构把查询拆到多个 BE 节点并行执行配合向量化执行引擎百万级数据的聚合查询通常能压到百毫秒级StreamLoad 和 Insert Into 支持实时写入数据进库即可查。MCP 解决的是接得顺、管得住这一层客户端通过工具发现接口自动拿到可用工具列表不用把表名、字段名写死认证和权限在协议层统一处理敏感库不至于裸奔。两者拼起来Agent 的实时分析链路就成型了数据实时进 DorisMCP 服务端把 Doris 的查询能力包装成标准工具Agent 按需调用拿到结果再组织语言回复用户。下面我从环境准备开始一步步把这套链路跑通。2. 前置准备TaoToken 侧与 Doris 侧各要什么先说结论这套链路里TaoToken 负责给 Agent 提供模型推理能力Doris 负责实时数据MCP 服务端是中间的粘合层。三者缺一不可但配置顺序建议先 Doris 再 MCP 再模型。Doris 侧你需要一个可访问的 FEFrontend和至少一个 BEBackend。本地测试用 Docker 起单机版最省事生产环境按官方文档做集群。关键连接参数有四个FE 的 MySQL 协议端口默认 9030HTTP 端口默认 8030BE 的 Web 端口默认 8040以及一个有查询权限的账号。建库建表用 MySQL 客户端连 9030 就行Doris 兼容 MySQL 协议这点对老手很友好。TaoToken 侧你需要一个 API Key用来让 Agent 调用模型。获取入口在控制台的 API Keys 页面登录后新建即可。模型对话能力可以直接在模型对话页验证长期跑编码或 Agent 任务的话Coding Plan 更划算按量计费适合先试水。接入文档里有各语言 SDK 的调用示例MCP 服务端里调模型的部分照着改就行。MCP 服务端本身建议用 Python 或 Node 写Python 生态里mcp官方 SDK 比较成熟。你需要准备一个能连 Doris 的数据库驱动pymysql或mysql-connector-python一个 MCP 服务端框架以及一个 HTTP 客户端用来调 TaoToken 的 API。目录结构建议这样分doris-mcp-server/ ├── server.py # MCP 服务端入口 ├── doris_client.py # Doris 连接与查询封装 ├── tools.py # 暴露给 Agent 的工具定义 └── config.yaml # 连接参数与密钥配置和密钥别写死在代码里用环境变量或配置文件加载后面排查问题也方便。Doris 账号建议单独建一个只读账号给 MCP 用别拿 root 直接上。3. 可复制配置MCP 服务端骨架与 Doris 连接参数这一节给可直接跑的代码。先装依赖pip install mcp pymysql pyyaml httpxconfig.yaml里放连接参数注意别把真实密码提交到仓库doris: host: 127.0.0.1 port: 9030 user: mcp_reader password: your_password database: analytics taotoken: base_url: https://taotoken.net/api api_key: sk-your-key model: claude-sonnetdoris_client.py封装连接和查询重点是加超时和行数上限防止 Agent 拉全表把 BE 打满import pymysql from pymysql.cursors import DictCursor class DorisClient: def __init__(self, cfg): self.cfg cfg def _conn(self): return pymysql.connect( hostself.cfg[host], portself.cfg[port], userself.cfg[user], passwordself.cfg[password], databaseself.cfg[database], charsetutf8mb4, cursorclassDictCursor, connect_timeout5, read_timeout30, ) def list_tables(self): with self._conn() as conn: with conn.cursor() as cur: cur.execute(SHOW TABLES) return [list(r.values())[0] for r in cur.fetchall()] def describe(self, table): with self._conn() as conn: with conn.cursor() as cur: cur.execute(fDESC {table}) return cur.fetchall() def query(self, sql, limit200): # 只允许 SELECT避免 Agent 误删数据 if not sql.strip().lower().startswith(select): raise ValueError(only SELECT is allowed) if limit not in sql.lower(): sql f{sql.rstrip(;)} LIMIT {limit} with self._conn() as conn: with conn.cursor() as cur: cur.execute(sql) return cur.fetchall()tools.py定义三个工具列表、表结构、执行查询。MCP 的工具描述要写清楚参数含义模型才知道怎么填from mcp.server import Server from mcp.types import Tool, TextContent import json def register_tools(server: Server, doris: DorisClient): server.list_tools() async def list_tools(): return [ Tool( namelist_tables, description列出 Doris 中所有可查询的表名, inputSchema{type: object, properties: {}}, ), Tool( namedescribe_table, description查看指定表的字段结构, inputSchema{ type: object, properties: {table: {type: string}}, required: [table], }, ), Tool( namerun_query, description执行只读 SELECT 查询返回 JSON 结果, inputSchema{ type: object, properties: {sql: {type: string}}, required: [sql], }, ), ] server.call_tool() async def call_tool(name, arguments): if name list_tables: data doris.list_tables() elif name describe_table: data doris.describe(arguments[table]) elif name run_query: data doris.query(arguments[sql]) else: raise ValueError(funknown tool: {name}) return [TextContent(typetext, textjson.dumps(data, ensure_asciiFalse, defaultstr))]server.py把上面拼起来用 stdio 传输启动import asyncio, yaml from mcp.server import Server from mcp.server.stdio import stdio_server from doris_client import DorisClient from tools import register_tools async def main(): cfg yaml.safe_load(open(config.yaml)) doris DorisClient(cfg[doris]) server Server(doris-mcp) register_tools(server, doris) async with stdio_server() as (r, w): await server.run(r, w, server.create_initialization_options()) if __name__ __main__: asyncio.run(main())跑起来就一行python server.py如果 Agent 客户端支持 HTTP 传输把stdio_server换成 SSE 或 streamable HTTP 即可工具注册逻辑不用动。Doris 连接参数里read_timeout别设太大Agent 等太久会超时重试反而放大压力。4. 验证请求从数据接入到 Agent 查询的完整闭环配置写完必须验证不然 Agent 调不通你都不知道卡在哪。分三步走。第一步确认 Doris 里有数据可查。用 MySQL 客户端连上去建个测试表CREATE TABLE IF NOT EXISTS sales ( dt DATE, region VARCHAR(32), product VARCHAR(64), amount DECIMAL(18,2) ) DUPLICATE KEY(dt, region, product) DISTRIBUTED BY HASH(product) BUCKETS 4 PROPERTIES (replication_num 1); INSERT INTO sales VALUES (2025-06-01,华东,A,1200), (2025-06-02,华东,B,800), (2025-06-03,华南,A,1500);第二步单独测 MCP 服务端的工具调用。用官方提供的 MCP Inspector 或者自己写个脚本直接调run_queryimport asyncio, json from mcp import ClientSession, StdioServerParameters from mcp.client.stdio import stdio_client async def test(): params StdioServerParameters(commandpython, args[server.py]) async with stdio_client(params) as (r, w): async with ClientSession(r, w) as session: await session.initialize() tools await session.list_tools() print(tools:, [t.name for t in tools.tools]) res await session.call_tool(run_query, { sql: SELECT region, SUM(amount) AS total FROM sales GROUP BY region ORDER BY total DESC }) print(res.content[0].text) asyncio.run(test())预期输出是工具列表[list_tables, describe_table, run_query]以及按区域汇总的销售额 JSON。如果这一步通了说明 Doris 到 MCP 的链路没问题。第三步接上 Agent 做端到端验证。在 Agent 客户端里配置 MCP 服务端然后问一句华东地区销售额最高的产品是哪个。Agent 会先调list_tables或describe_table摸清结构再生成 SQL 调run_query最后把结果组织成自然语言。实测下来从提问到返回结果整条链路在本地环境通常几百毫秒内完成瓶颈往往在模型推理而不是 Doris 查询。验证时重点看两件事一是 Agent 生成的 SQL 有没有越界比如没加 LIMIT二是返回结果有没有被正确解析。前者靠doris_client.py里的白名单和自动补 LIMIT 兜底后者看 MCP 返回的 JSON 是否合法。5. 本篇常见错排查报错一pymysql.err.OperationalError: (2003, Cant connect to MySQL server)八成是 FE 的 9030 端口没通。先telnet 127.0.0.1 9030确认再检查 Doris 的fe.conf里query_port配置。Docker 部署的话注意端口映射别只映射 8030 忘了 9030。报错二Access denied for user mcp_reader账号权限没给够。Doris 里执行GRANT SELECT_PRIV ON analytics.* TO mcp_reader%;然后FLUSH PRIVILEGES;。注意 Doris 的权限模型和 MySQL 略有差异库级权限用SELECT_PRIV。报错三MCP 客户端连不上服务端日志显示Connection closed多半是server.py启动就崩了。单独跑python server.py看报错常见的是config.yaml路径不对或 YAML 缩进错误。stdio 模式下服务端不能往 stdout 打日志否则会污染协议流日志一律走 stderr。报错四Agent 生成的 SQL 报Syntax error near ...Doris 的 SQL 方言和 MySQL 有差异比如日期函数、窗口函数写法。让 Agent 先调describe_table拿到字段类型再生成 SQL命中率会高很多。也可以在工具描述里加一句请使用 Doris 兼容语法。报错五查询很慢BE CPU 打满检查是不是没加分区过滤或 LIMIT。Doris 的 MPP 架构对全表扫描不友好Agent 生成的 SQL 如果没带 WHERE 条件很容易扫全表。在run_query里强制要求带时间范围或 LIMIT能挡掉大部分慢查询。报错六TaoToken 侧返回 401API Key 没配对或者 base_url 写成了带 UTM 的地址。注意 API 调用地址是https://taotoken.net/api不要带查询参数。Key 建议放环境变量别硬编码进config.yaml提交。6. 把链路跑稳之后往哪走跑通上面这套之后你会发现真正的难点不在 Doris 也不在 MCP而在Agent 怎么生成靠谱的 SQL。我的经验是工具描述写得越具体模型瞎猜的概率越低。比如run_query的描述里可以补一句表 sales 的 dt 字段是分区键查询请带上 dt 范围比让模型自己describe一遍再猜要快。另一个实用技巧是给 MCP 服务端加一层查询缓存。Agent 在同一个会话里经常重复问相似问题把(sql, 结果)缓存几十秒能明显降低 Doris 压力。缓存 key 用 SQL 的哈希注意别缓存带NOW()这类非确定性函数的查询。如果你要长期跑 Agent 任务建议把模型调用切到 Coding Plan按量计费在频繁工具调用的场景下更可控。接入细节看接入文档里面有 MCP 场景的完整示例。模型对话页可以先用几轮对话验证工具调用是否符合预期再上生产。最后提醒一句Doris 的实时写入和 MCP 的工具发现是两套独立机制别指望 MCP 帮你管数据同步。数据接入该用 StreamLoad 就用 StreamLoad该用 Routine Load 就用 Routine LoadMCP 只负责把已经能查的数据暴露给 Agent。职责分清链路才稳。
返回列表