1. 数据湖 ETL 的老问题:文件堆成山,AI 却看不见
数据湖里最尴尬的场景不是没有数据,而是数据就在对象存储里躺着,Parquet 文件按天分区堆了几十万个,但想让 AI 帮你写一段 ETL 把最近三个月的订单去重合并,它连文件路径都摸不到。传统做法是工程师先手写 Spark 任务或者一段长 SQL,跑通了再交给调度器,源端字段一改,整条 Pipeline 就崩,排查成本极高。
MCP(Model Context Protocol)解决的正是这个断层:它把 DuckDB 这类嵌入式分析引擎和 Parquet 文件系统封装成 AI 可以直接调用的工具与资源,让模型自己探测 Schema、生成转换 SQL、执行并校验结果。DuckDB 在这里的角色很关键,它单进程启动、毫秒级响应、对 Parquet 的向量化读取性能顶尖,非常适合做数据湖的“探索层”和“即时 ETL 验证层”,而不是每次都拉起沉重的分布式集群。
这篇面向数据湖场景,交付一个可复制的 MCP Server 骨架,把 DuckDB + Parquet 的分析能力暴露给 AI,并用 TaoToken 统一 Key 打通模型调用通道,让 AI 自动完成从 Schema 探测到 ETL 执行的完整编排。适合正在做数据平台、想用 AI 降低 ETL 开发门槛的工程师跟做。
2. TaoToken 前置:统一 Key 与 API 通道准备
MCP Server 本身只负责“执行”,真正做编排决策的是背后的模型。你需要一个稳定的模型调用通道,TaoToken 在这里提供统一 Key,把模型对话、编码类模型、Agent 编排都收敛到一套凭证上,省去在多个供应商之间切换配置的麻烦。
接入分三步走。第一步,到官网注册并进入控制台,地址是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_medium=csdn&utm_campaign=rewrite&utm_content= ,登录后在控制台里创建 API Key,入口在 https://taotoken.net/console?utm_source=taotoken_aicg_blog_end&utm_content=console&utm_campaign=rewrite 。第二步,如果你要长期跑编码类或 Agent 类任务,建议看一下 Coding Plan,路径是 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite ,它更适合高频调用场景。第三步,Key 管理页面在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite ,生成后复制保存,后面写进 MCP Server 的环境变量。
API 基础地址统一用 https://taotoken.net/api ,注意这个地址不带 UTM 参数,直接作为 base_url 使用。如果你用的是 Claude Code 这类工具做 Agent 编排,可以参考 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite 里的接入说明,把 base_url 和 Key 填进去即可。想先验证模型是否通,可以直接在模型对话页 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_content=model-chat&utm_campaign=rewrite 发一条消息测试。
注意:Key 只放在服务端环境变量里,不要写进前端代码或提交到 Git 仓库。MCP Server 通过 process.env 读取,避免泄露。
3. 可复制配置:MCP Server 骨架与 DuckDB 集成
先建项目目录并装依赖。DuckDB 的 Node 驱动加上 MCP SDK 就够了,不需要额外装 S3 客户端,DuckDB 的 httpfs 扩展会处理远程读取。
mkdir mcp-datalake-server && cd mcp-datalake-server npm init -y npm install @modelcontextprotocol/sdk duckdb npm install -D typescript @types/node tsx npx tsc --init接着写核心 Server。它暴露两个工具:一个探测 Parquet Schema,一个执行 AI 生成的 ETL SQL。DuckDB 用内存模式作为临时转换站,通过 httpfs 直接读远程 Parquet,做到零拷贝。
import { Server } from "@modelcontextprotocol/sdk/server/index.js"; import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js"; import { ListToolsRequestSchema, CallToolRequestSchema, } from "@modelcontextprotocol/sdk/types.js"; import duckdb from "duckdb"; const server = new Server( { name: "datalake-navigator", version: "1.0.0" }, { capabilities: { tools: {}, resources: {} } } ); const db = new duckdb.Database(":memory:"); const con = db.connect(); // 加载 httpfs 以支持远程 Parquet 读取 con.run(`INSTALL httpfs; LOAD httpfs;`); con.run(`SET s3_region='us-east-1';`); server.setRequestHandler(ListToolsRequestSchema, async () => ({ tools: [ { name: "explore_table_schema", description: "探测远程 Parquet 文件的表结构和基础元数据。调用前请确认路径包含分区字段。", inputSchema: { type: "object", properties: { file_path: { type: "string", description: "Parquet 路径,如 's3://bucket/data/*.parquet'", }, }, required: ["file_path"], }, }, { name: "execute_autonomous_etl", description: "执行 AI 生成的 DuckDB SQL 完成 ETL。SQL 必须包含分区字段过滤,否则会被拒绝。", inputSchema: { type: "object", properties: { sql: { type: "string", description: "DuckDB 语法 SQL" }, operation_type: { type: "string", enum: ["TRANSFORM", "AGGREGATE", "CLEAN"], }, }, required: ["sql"], }, }, ], })); server.setRequestHandler(CallToolRequestSchema, async (request) => { const { name, arguments: args } = request.params; if (name === "explore_table_schema") { const path = args?.file_path as string; const query = `DESCRIBE SELECT * FROM read_parquet('${path}') LIMIT 0;`; return new Promise((resolve) => { con.all(query, (err, res) => { if (err) { resolve({ content: [{ type: "text", text: `探测失败: ${err.message}` }], isError: true, }); return; } resolve({ content: [ { type: "text", text: `Schema 结果:\n${JSON.stringify(res, null, 2)}` }, ], }); }); }); } if (name === "execute_autonomous_etl") { const sql = args?.sql as string; // 成本熔断:无分区过滤直接拒绝,防止全量扫描 if (!/dt\s*=/.test(sql) && !/dt\s+BETWEEN/i.test(sql)) { return { content: [ { type: "text", text: "拒绝执行:SQL 缺少 dt 分区过滤条件。" }, ], isError: true, }; } return new Promise((resolve) => { con.all(sql, (err, res) => { if (err) { resolve({ content: [{ type: "text", text: `ETL 失败: ${err.message}` }], isError: true, }); return; } resolve({ content: [ { type: "text", text: `执行结果预览:\n${JSON.stringify(res?.slice(0, 20), null, 2)}`, }, ], }); }); }); } throw new Error("Tool not found"); }); const transport = new StdioServerTransport(); await server.connect(transport);配置部分,MCP 客户端通过 config.toml 或 settings.json 挂载这个 Server。以常见的 MCP 客户端配置为例,config.toml 写法如下:
[mcp_servers.datalake] command = "npx" args = ["tsx", "/path/to/mcp-datalake-server/index.ts"] env = { TAOTOKEN_API_KEY = "你的Key", TAOTOKEN_BASE_URL = "https://taotoken.net/api" }如果你用的是 JSON 配置的客户端,settings.json 等价写法:
{ "mcpServers": { "datalake": { "command": "npx", "args": ["tsx", "/path/to/mcp-datalake-server/index.ts"], "env": { "TAOTOKEN_API_KEY": "你的Key", "TAOTOKEN_BASE_URL": "https://taotoken.net/api" } } } }提示:路径用绝对路径,避免客户端工作目录不同导致找不到入口文件。Key 通过 env 注入,Server 内部用 process.env.TAOTOKEN_API_KEY 读取。
4. 验证请求:用 Parquet 样本跑通 AI 自动 ETL
准备一份样本 Parquet 数据,本地或对象存储都行。这里用本地文件模拟,路径换成你的实际路径即可。假设文件是 orders.parquet,字段有 order_id、user_id、amount、dt。
第一步,让 AI 探测 Schema。在 MCP 客户端里发指令:“用 datalake 工具探测 /data/orders.parquet 的结构”。AI 会调用 explore_table_schema,返回字段列表和类型。你会看到类似 order_id VARCHAR、amount DOUBLE、dt VARCHAR 的结果。
第二步,让 AI 生成 ETL。指令:“按 dt 过滤 2024-01-01 到 2024-01-31,对 user_id 聚合 amount 求和,去重 order_id”。AI 会生成一段 DuckDB SQL,类似:
SELECT user_id, SUM(amount) AS total_amount FROM ( SELECT DISTINCT order_id, user_id, amount FROM read_parquet('/data/orders.parquet') WHERE dt BETWEEN '2024-01-01' AND '2024-01-31' ) GROUP BY user_id ORDER BY total_amount DESC;第三步,AI 调用 execute_autonomous_etl 执行。Server 先检查 SQL 是否含 dt 过滤,通过后执行,返回前 20 条预览。你会在客户端看到聚合结果,比如 user_id 为 u1001 的 total_amount 是 3820.5。整个过程 AI 自己完成探测、生成、执行、校验,你只需要描述意图。
如果想验证模型通道是否正常,可以先用模型对话页发一条简单消息确认 Key 有效,地址是 https://taotoken.net/?utm_source=taotoken_aicg_blog_end&utm_content=model-chat&utm_campaign=rewrite 。确认后再跑 MCP 编排,避免把通道问题和 Server 问题混在一起排查。
5. 本篇常见错排查
报错一:IO Error: Extension "httpfs" not found。DuckDB 首次加载 httpfs 需要联网下载扩展。如果环境无外网,提前在有网机器上执行 INSTALL httpfs,把生成的扩展文件拷贝到目标机器的 .duckdb/extensions 目录。或者改用本地文件路径测试,先跑通逻辑再上远程。
报错二:MCP 客户端连不上 Server,无任何输出。多半是入口路径不对或 tsx 未安装。检查 config.toml 里的 args 是否指向真实存在的 index.ts,并在项目目录手动执行 npx tsx index.ts 看是否报错。Stdio 模式下 Server 不能往 stdout 打日志,所有调试信息走 stderr,否则会污染 MCP 协议消息。
报错三:SQL 被拒绝,提示缺少 dt 分区过滤。这是成本熔断逻辑在起作用。如果你的表分区字段不叫 dt,改一下正则匹配的字段名。这个校验的目的是防止 AI 写出 SELECT * 全量扫描,对象存储的流量费和计算费会瞬间上涨,带成本感知的 API 是数据湖治理的关键。
报错四:Schema 探测返回空。检查 Parquet 路径是否用了通配符,read_parquet 支持 's3://bucket/data/*.parquet' 这种写法。如果路径正确但返回空,可能是 S3 凭证没配,需要在 con.run 里补 SET s3_access_key_id 和 s3_secret_access_key,或者用环境变量注入。
报错五:AI 生成的 SQL 字段名对不上。这是语义幻觉,文件里字段可能是 col_01 这种无意义命名。对策是在探测阶段让 AI 同时计算 distinct_count 和样本值,辅助它猜测字段含义。更彻底的做法是维护一份 metadata.json,把 col_01 映射为 customer_id,在 Resource 里暴露给 AI。
6. 长期编码与 Agent 编排的通道选择
如果你只是偶尔跑一次 ETL 验证,按需调用模型对话就够了。但如果要把这套 MCP Server 接进日常数据开发流程,让 AI 持续做 Schema 探测、SQL 生成、结果校验,调用频率会很高,这时候用 Coding Plan 更划算,入口在 https://taotoken.net/coding-plan?utm_source=taotoken_aicg_blog_end&utm_content=coding-plan&utm_campaign=rewrite 。它适合长期编码和 Agent 类任务,统一 Key 管理也省心。
接入文档在 https://taotoken.net/doc?utm_source=taotoken_aicg_blog_end&utm_content=doc&utm_campaign=rewrite ,里面有 base_url 配置和常见客户端接入示例。API Key 管理在 https://taotoken.net/api-keys?utm_source=taotoken_aicg_blog_end&utm_content=api-keys&utm_campaign=rewrite ,随时可以轮换。API 基础地址固定用 https://taotoken.net/api ,不带 UTM 参数。
实测下来,这套组合最实用的地方在于把 ETL 的开发周期从“写代码-调试-调度”压缩到“描述意图-校验结果”。DuckDB 负责零拷贝探索,MCP 负责把能力暴露给 AI,TaoToken 负责统一模型通道,三者各司其职。你可以先从本地 Parquet 样本跑通,再逐步换成对象存储路径,最后把分区过滤和成本熔断加上,就是一套能进生产环境的数据湖智能 ETL 骨架。