火山引擎金融Agent数据底座实战:A股行情接入、豆包Tool Call与扣子工作流

一、Agent 会猜价格,是因为它没有数据底座

早期做A股策略,我用Python拉通了日K线,回测曲线挺好看,以为数据层搞定了。实盘上线后问题一个接一个:除权日跳空被策略当成下跌信号,两天后才反应过来是K线没做复权;想估滑点,发现手里只有K线,没有盘口,只能拍一个固定值。

后来我把同样的策略交给 Agent 跑,问题更严重:问它"现在某某标的价格多少",它会用训练数据里的历史价格回答,而不是去调接口取最新值。 问它"这只票最近资金流怎么样",它会用"一般来说主力资金……"这种没有事实依据的表述。

这不是模型能力问题,是Agent 的数据底座没搭好。模型负责推理,数据负责事实。让模型用记忆猜价格,是把分析过程变成幻觉生成过程。

这篇文章要讲一件事:在火山引擎上,如何给金融 Agent 搭一套能取到带时间戳、带标的、带字段的事实的数据底座。 从行情接入、数据管道、工具封装,到火山方舟 Tool Call 和扣子工作流编排,全程具体代码,不给概念图。


目录


二、Agent 数据底座五层架构

┌──────────────────────────────────────────────────────────────────┐
│                      火山引擎 VPC                                  │
│                                                                    │
│   ┌────────────┐   ┌────────────┐   ┌────────────┐               │
│   │ 接入层      │   │ 传输层      │   │ 存储层      │               │
│   │ REST/WS/MCP│──▶│ Kafka 版   │──▶│ TOS + veDB │               │
│   └─────┬──────┘   └────────────┘   └──────┬─────┘               │
│         │                                    │                     │
│         │ 行情源                             │                     │
│         ▼                                    ▼                     │
│   ┌────────────┐                    ┌────────────────┐            │
│   │ TickDB     │                    │ 工具层          │            │
│   │ REST + WS  │                    │ MCP / Tool Call│            │
│   └────────────┘                    └───────┬────────┘            │
│                                              │                     │
│                                              ▼                     │
│                                     ┌─────────────────┐           │
│                                     │ Agent 层         │           │
│                                     │ 火山方舟 + 扣子  │           │
│                                     └─────────────────┘           │
│                                                                    │
│   ┌────────────────────────────────────────────────────────┐     │
│   │ 监控:TLS 日志 + 火山引擎可观测                        │     │
│   └────────────────────────────────────────────────────────┘     │
└──────────────────────────────────────────────────────────────────┘

五层各解决一个问题:

解决火山引擎产品关键约束
接入拿数据veFaaS / ECS / 轻量服务器密钥管理、重试、幂等
传输削峰解耦消息队列 Kafka 版分区键、消息顺序
存储冷热分离TOS + veDB冷存归档、热存索引
工具让 Agent 能调MCP Server / Tool Call Schema先取事实、再推理
Agent编排推理火山方舟 + 扣子事实先于推理

这五层里,工具层是 Agent 场景的关键。 少了工具层,前四层只是数据管道;有了工具层,数据才能被 Agent 调用。


三、接入层:行情 API 的 REST + WebSocket + MCP

3.1 行情接口的基础结构

先把鉴权和响应结构固定下来,后面 veFaaS、ECS、VKE 里所有代码复用同一套封装。

REST 端点是 https://api.tickdb.ai/v1,认证 Header 是 X-API-Key,响应统一为 {code, message, data} 包装,code=0 表示成功。不同端点的 data 结构不一样——ticker 是数组,kline 是带 klines 的对象,估值是 metrics 嵌套,不能假设结构一致。

# 测试环境:API 行为实测于 Python 3.14.2 / macOS 15.5 / 2026-09-16
# 火山引擎 veFaaS Python 3.11 运行时兼容,需在函数配置里注入 TICKDB_API_KEY

import os
import requests

BASE_URL = "https://api.tickdb.ai/v1"
API_KEY = os.getenv("TICKDB_API_KEY")

def _headers() -> dict:
    if not API_KEY:
        raise RuntimeError("请在 veFaaS 环境变量中配置 TICKDB_API_KEY")
    return {"X-API-Key": API_KEY}

def unwrap(resp: requests.Response) -> dict:
    """统一解包 {code, message, data},认证错误先于业务错误被捕获"""
    resp.raise_for_status()
    payload = resp.json()
    if payload.get("code") != 0:
        raise RuntimeError(f"业务错误: code={payload.get('code')} message={payload.get('message')}")
    return payload["data"]

def get_ticker(symbol: str) -> dict:
    """实时快照,data 是数组"""
    url = f"{BASE_URL}/market/ticker"
    resp = requests.get(url, headers=_headers(), params={"symbols": symbol}, timeout=10)
    data = unwrap(resp)
    if not data:
        raise RuntimeError(f"未找到标的: {symbol}")
    return data[0]

def get_kline(symbol: str, beg_ms: int, end_ms: int,
              interval: str = "1d", adjust: str = "forward", limit: int = 1000) -> list:
    """
    历史K线
    注意:时间参数是 start_time/end_time 毫秒值,不是 beg_day/end_day。
    beg_day/end_day 返回 200 但时间窗不生效——200 不等于参数生效。
    """
    url = f"{BASE_URL}/market/kline"
    params = {
        "symbol": symbol, "interval": interval, "adjust": adjust,
        "start_time": beg_ms, "end_time": end_ms, "limit": limit,
    }
    resp = requests.get(url, headers=_headers(), params=params, timeout=15)
    data = unwrap(resp)
    return data.get("klines", [])

3.2 WebSocket 常驻订阅

快照是定时拉,但逐笔和盘口必须是流式订阅。veFaaS 不适合长连接(执行时间上限、冷启动会断开),这部分放到 ECS 或轻量应用服务器。

WebSocket 连接 URL 是 wss://api.tickdb.ai/v1/realtime?api_key=...,实测订阅消息格式是 {"cmd":"subscribe","data":{"channel":"ticker","symbols":["..."]}}不是 action/channel/symbol 平铺格式——我一开始按直觉写,服务端返回 code=2001 Unknown command,排查了半小时。

# 部署在 ECS 上,systemd 管理常驻进程
import os
import ssl
import json
import certifi
import asyncio
import websockets
from kafka import KafkaProducer

api_key = os.getenv("TICKDB_API_KEY")

async def stream_to_kafka(symbols: list):
    url = f"wss://api.tickdb.ai/v1/realtime?api_key={api_key}"
    ssl_ctx = ssl.create_default_context(cafile=certifi.where())
    
    producer = KafkaProducer(
        bootstrap_servers=os.environ["KAFKA_BROKER"].split(","),
        value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode("utf-8"),
        key_serializer=lambda k: k.encode("utf-8") if k else None,
        acks="all", retries=3, linger_ms=50,
    )
    
    retry_backoff = 1
    while True:
        try:
            async with websockets.connect(url, ssl=ssl_ctx, ping_interval=20) as ws:
                sub = {"cmd": "subscribe", "data": {"channel": "ticker", "symbols": symbols}}
                await ws.send(json.dumps(sub))
                print(f"已订阅: {sub}")
                retry_backoff = 1
                
                async for message in ws:
                    # 行情推送字段待交易时段补测,这里按原样转发
                    producer.send(
                        os.environ["KAFKA_TOPIC_WS"],
                        key=None,
                        value={"source": "ws", "raw": message},
                    )
        except websockets.exceptions.ConnectionClosedError as e:
            if e.code == 1008:
                # 密钥过期必须人工介入,不能自动重试
                print("密钥过期,停止重连,请刷新 API Key")
                raise SystemExit(1)
            print(f"连接异常关闭 code={e.code}{retry_backoff}s 后重连")
            await asyncio.sleep(retry_backoff)
            retry_backoff = min(retry_backoff * 2, 60)
        except Exception as e:
            print(f"未知异常: {e}{retry_backoff}s 后重连")
            await asyncio.sleep(retry_backoff)
            retry_backoff = min(retry_backoff * 2, 60)

if __name__ == "__main__":
    asyncio.run(stream_to_kafka(["688256.SH", "000001.SZ"]))

两个关键点

  1. close code 1008 = 密钥过期,必须停止重连并人工介入。普通断线才用指数退避重连。所有断线一视同仁地疯狂重试,账号可能被锁。
  2. 显式使用 certifi 提供的 CA。实测在 macOS 上用系统证书链连接失败,改用 certifi 后成功。ECS 上同理。

3.3 veFaaS 定时快照采集

veFaaS 是最省的方案:定时触发,只跑几十毫秒,拉快照,推 Kafka,退出。无需常驻进程。

veFaaS 函数配置:

配置项
运行环境Python 3.11
内存256 MB
超时15 秒
触发器定时触发器,1 分钟一次
VPC与消息队列 Kafka 版同一 VPC
环境变量TICKDB_API_KEY / KAFKA_BROKER / KAFKA_TOPIC
# veFaaS 入口函数:定时拉取快照并推 Kafka
import json
import os
from kafka import KafkaProducer

SYMBOLS = os.getenv("SYMBOLS", "688256.SH,000001.SZ").split(",")

def handler(event, context):
    producer = KafkaProducer(
        bootstrap_servers=os.environ["KAFKA_BROKER"].split(","),
        value_serializer=lambda v: json.dumps(v, ensure_ascii=False).encode("utf-8"),
        key_serializer=lambda k: k.encode("utf-8") if k else None,
        acks="all",              # 保证写入确认
        retries=3,               # 网络抖动重试
        enable_idempotence=True, # 幂等生产,避免重复消息
        linger_ms=50,
    )
    
    results = []
    for sym in SYMBOLS:
        try:
            tick = get_ticker(sym)
            # 分区键用 symbol,保证同一标的的消息按顺序进入同一分区
            producer.send(
                os.environ["KAFKA_TOPIC"],
                key=sym,
                value={"source": "ticker", "symbol": sym, "payload": tick},
            )
            results.append({"symbol": sym, "ok": True})
        except Exception as e:
            results.append({"symbol": sym, "ok": False, "error": str(e)})
    
    producer.flush(timeout=10)
    print(json.dumps({"count": len(results), "detail": results}, ensure_ascii=False))
    return {"statusCode": 200, "body": json.dumps(results, ensure_ascii=False)}

部署要点

  • veFaaS 默认公网出口,访问 Kafka 版必须配置 VPC 与子网,否则连接不通。
  • 环境变量在"函数配置 → 环境变量"里注入,不写进代码。
  • 定时触发器用 Cron 表达式,A股交易时段 9:15–15:00 触发。

四、传输层:消息队列 Kafka 版分区与幂等

Kafka 在管道里做两件事:削峰解耦。采集端按分钟推,Agent 按需拉,中间必须有缓冲。

4.1 Topic 设计

Topic用途分区数保留
a-stock-ticker快照624 小时symbol
a-stock-wsWebSocket 原始推送612 小时symbol
a-stock-dlq死信队列17 天

4.2 用 symbol 做分区键

producer.send("a-stock-ticker", key=symbol, value=payload)

同一个 symbol 的消息进入同一个分区,保证消费端按序处理。 不用 key,Kafka 会轮询分区,同一标的的快照可能乱序到达。

4.3 幂等生产

producer = KafkaProducer(
    bootstrap_servers=BROKERS,
    acks="all",              # 所有 ISR 副本确认才算写入成功
    retries=3,
    enable_idempotence=True, # 幂等生产
    max_in_flight_requests_per_connection=1,
)

acks=all + enable_idempotence=True 是金融数据场景的底线。丢一条 K 线,回测结果就偏了;重复一条,因子计算就出 NaN。

4.4 消息体格式

统一信封,方便消费端按 source 分发:

{
  "source": "ticker",
  "symbol": "688256.SH",
  "ts_ingest": 1754118000123,
  "payload": {
    "symbol": "688256.SH",
    "last_price": "1086.58",
    "open": "1055.00",
    "prev_close": "1053.98",
    "timestamp": 1754118000000
  }
}

ts_ingest 是入队时间戳,payload.timestamp 是数据本身的时间戳。这两个不能混——延迟统计和去重都靠它们。


五、存储层:TOS 冷数据 + veDB 热数据

冷热分离是必须的。把 10 年 K 线全塞进 MySQL,查询会慢到没法用;把最近 5 分钟快照全写 TOS,每次读都要等几秒。

5.1 TOS 归档历史 K 线

历史 K 线按 symbol/interval/year/month 分层存储,Parquet 格式,按天分区。

# VKE Job 或 veFaaS 定时任务:拉取昨日K线归档到 TOS
import io
import os
import pandas as pd
from datetime import datetime, timedelta, timezone
import tos

def archive_daily_kline(symbol: str, trade_date: datetime):
    beg_ms = int((trade_date.replace(hour=0, minute=0, second=0, microsecond=0)
                  .replace(tzinfo=timezone.utc)).timestamp() * 1000)
    end_ms = beg_ms + 24 * 3600 * 1000
    
    klines = get_kline(symbol, beg_ms, end_ms, interval="1m")
    if not klines:
        return None
    
    df = pd.DataFrame(klines)
    df["symbol"] = symbol
    df["trade_date"] = trade_date.strftime("%Y-%m-%d")
    
    buf = io.BytesIO()
    df.to_parquet(buf, index=False, compression="snappy")
    buf.seek(0)
    
    # TOS 分层路径
    key = (f"a-stock/kline/{symbol}/"
           f"year={trade_date.year}/month={trade_date.month:02d}/"
           f"{trade_date.strftime('%Y-%m-%d')}.parquet")
    
    client = tos.TosClientV2(
        ak=os.environ["TOS_ACCESS_KEY"],
        sk=os.environ["TOS_SECRET_KEY"],
        endpoint=os.environ["TOS_ENDPOINT"],
        region=os.environ["TOS_REGION"],
    )
    client.put_object(bucket=os.environ["TOS_BUCKET"], key=key, content=buf)
    print(f"归档完成: tos://{os.environ['TOS_BUCKET']}/{key}")
    return key

为什么选 Parquet + Snappy

  • 列式存储,因子计算时只读需要的列,比 CSV 快 5–10 倍
  • Snappy 压缩比约 3:1,10 年 1 分钟 K 线压缩后不到 20 GB
  • TOS 标准存储单价低,长期归档再转低频,进一步降本

5.2 veDB 存热数据与元信息

veDB 只存两类

  1. 最近 N 天的快照和逐笔,供实时看板和 Agent 查询
  2. 标的基础信息、交易日历、复权因子等元数据,量小但查询频繁
CREATE TABLE ticker_snapshot (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  symbol VARCHAR(32) NOT NULL,
  last_price DECIMAL(18,4),
  open_price DECIMAL(18,4),
  prev_close DECIMAL(18,4),
  volume BIGINT,
  data_ts BIGINT NOT NULL,
  ingest_ts BIGINT NOT NULL,
  KEY idx_symbol_ts (symbol, data_ts),
  KEY idx_ingest (ingest_ts)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

CREATE TABLE trade_calendar (
  market VARCHAR(8) NOT NULL,
  trade_date DATE NOT NULL,
  is_half_day TINYINT DEFAULT 0,
  PRIMARY KEY (market, trade_date)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

索引要点(symbol, data_ts) 组合索引用于"查某标的某时间段"的查询。ingest_ts 单独索引用于延迟统计。

5.3 冷热切换策略

热数据(veDB):最近 7 天
    ↓ 每日 02:00 归档任务
温数据(TOS 标准):7 天 - 1 年
    ↓ 生命周期策略
冷数据(TOS 低频):1 年 - 3 年
    ↓ 生命周期策略
归档数据(TOS 归档):3 年以上

六、工具层:把行情接口封装成 Agent Tool

这一层是 Agent 场景的核心。 前四层是数据管道,工具层是让 Agent 能"调用"数据的接口。没有工具层,Agent 只能靠模型记忆;有了工具层,Agent 才能取到带时间戳的事实。

6.1 MCP Server:标准化的 Agent 工具接口

MCP(Model Context Protocol)是 Agent 调用外部工具的标准协议。把行情接口封装成 MCP Server,Agent 就能通过标准协议调用。

# mcp_server.py —— 用 mcp SDK 暴露行情工具
# pip install mcp requests
import os
import json
import requests
from mcp.server import Server
from mcp.server.stdio import stdio_server
from mcp.types import Tool, TextContent

app = Server("tickdb-market-data")
BASE_URL = "https://api.tickdb.ai/v1"

def _headers():
    return {"X-API-Key": os.environ["TICKDB_API_KEY"]}

def _unwrap(resp):
    resp.raise_for_status()
    p = resp.json()
    if p.get("code") != 0:
        raise RuntimeError(f"业务错误: {p.get('message')}")
    return p["data"]

@app.list_tools()
async def list_tools() -> list[Tool]:
    """声明 Agent 可以调用的工具"""
    return [
        Tool(
            name="get_realtime_ticker",
            description="获取A股/美股/港股实时快照。返回最新价、开盘价、昨收价、成交量、时间戳。用于回答'现在多少钱'这类需要实时事实的问题。",
            inputSchema={
                "type": "object",
                "properties": {
                    "symbol": {"type": "string", "description": "标的代码,如 688256.SH"}
                },
                "required": ["symbol"],
            },
        ),
        Tool(
            name="get_kline",
            description="获取历史K线。时间参数用毫秒时间戳。返回 OHLCV 数组。用于回测、因子计算、趋势分析。",
            inputSchema={
                "type": "object",
                "properties": {
                    "symbol": {"type": "string"},
                    "beg_ms": {"type": "integer", "description": "起始毫秒时间戳"},
                    "end_ms": {"type": "integer", "description": "结束毫秒时间戳"},
                    "interval": {"type": "string", "default": "1d"},
                    "adjust": {"type": "string", "default": "forward"},
                },
                "required": ["symbol", "beg_ms", "end_ms"],
            },
        ),
        Tool(
            name="get_valuation",
            description="获取估值快照。返回 PE/PB/PS/DvdYld。用于估值分位分析。",
            inputSchema={
                "type": "object",
                "properties": {"symbol": {"type": "string"}},
                "required": ["symbol"],
            },
        ),
    ]

@app.call_tool()
async def call_tool(name: str, arguments: dict) -> list[TextContent]:
    """执行 Agent 的工具调用"""
    if name == "get_realtime_ticker":
        resp = requests.get(f"{BASE_URL}/market/ticker",
                            headers=_headers(),
                            params={"symbols": arguments["symbol"]},
                            timeout=10)
        data = _unwrap(resp)
        if not data:
            return [TextContent(type="text", text=f"未找到标的: {arguments['symbol']}")]
        return [TextContent(type="text", text=json.dumps(data[0], ensure_ascii=False))]
    
    elif name == "get_kline":
        resp = requests.get(f"{BASE_URL}/market/kline",
                            headers=_headers(),
                            params={
                                "symbol": arguments["symbol"],
                                "interval": arguments.get("interval", "1d"),
                                "adjust": arguments.get("adjust", "forward"),
                                "start_time": arguments["beg_ms"],
                                "end_time": arguments["end_ms"],
                                "limit": 1000,
                            },
                            timeout=15)
        data = _unwrap(resp)
        klines = data.get("klines", [])
        return [TextContent(type="text", text=json.dumps(klines, ensure_ascii=False))]
    
    elif name == "get_valuation":
        resp = requests.get(f"{BASE_URL}/fundamentals/valuation/latest",
                            headers=_headers(),
                            params={"symbol": arguments["symbol"]},
                            timeout=10)
        data = _unwrap(resp)
        # 实测路径:data.metrics["PE"]["value"]
        metrics = data.get("metrics", {})
        result = {
            "symbol": arguments["symbol"],
            "PE": metrics.get("PE", {}).get("value"),
            "PB": metrics.get("PB", {}).get("value"),
            "PS": metrics.get("PS", {}).get("value"),
            "DvdYld": metrics.get("DvdYld", {}).get("value"),
        }
        return [TextContent(type="text", text=json.dumps(result, ensure_ascii=False))]
    
    return [TextContent(type="text", text=f"未知工具: {name}")]

async def main():
    async with stdio_server() as (read, write):
        await app.run(read, write, app.create_initialization_options())

if __name__ == "__main__":
    import asyncio
    asyncio.run(main())

这个 MCP Server 的设计要点

  1. 工具描述写清"用于什么场景",让模型知道什么时候调用。比如"用于回答'现在多少钱'这类需要实时事实的问题"——这是给模型看的提示。
  2. 参数用 inputSchema 明确类型和必填项,减少模型传参错误。
  3. 返回值统一 JSON 字符串,让模型能解析。
  4. 单次调用边界:任何一次调用成功只证明该时间、样本、参数下的接口行为。

6.2 在火山方舟里注册 MCP 工具

在火山方舟控制台,创建应用后:

  1. 进入「应用配置」→「工具」→「添加 MCP Server」
  2. 填入 MCP Server 的启动命令(python mcp_server.py
  3. 在环境变量里注入 TICKDB_API_KEY
  4. 选择要暴露给 Agent 的工具

注册完成后,Agent 在对话中就能自动调用这些工具。

6.3 传统 Tool Call Schema 方案

如果不用 MCP,也可以直接用火山方舟的 Function Calling 格式:

TOOLS = [
    {
        "type": "function",
        "function": {
            "name": "get_realtime_ticker",
            "description": "获取A股实时快照,用于回答'现在多少钱'这类需要实时事实的问题。",
            "parameters": {
                "type": "object",
                "properties": {
                    "symbol": {"type": "string", "description": "标的代码,如 688256.SH"}
                },
                "required": ["symbol"],
            },
        },
    },
    # ... 其他工具
]

两种方案的选择:MCP 适合跨平台复用,一次封装,火山方舟、扣子、Claude 都能调;Function Calling 适合单平台深度集成,配置更直接。


七、Agent 层:火山方舟 Tool Call 与扣子工作流

7.1 火山方舟:用豆包做金融推理

火山方舟是字节的大模型服务平台,豆包系列模型在这里。调用豆包时,关键是让模型先调工具取事实,再基于事实推理

# pip install volcengine-python-sdk
import os
import json
from volcenginesdkarkruntime import Ark

client = Ark(api_key=os.environ["ARK_API_KEY"])

def ask_with_tools(question: str, model: str = "doubao-pro-32k"):
    """带工具调用的对话,模型先取事实再回答"""
    messages = [
        {
            "role": "system",
            "content": (
                "你是金融数据分析助手。回答任何涉及具体价格、成交、估值的问题前,"
                "必须先调用工具获取实时数据。不要用训练数据里的历史价格回答。"
                "如果工具返回空数组,明确告知用户'未找到该标的',不要编造。"
            ),
        },
        {"role": "user", "content": question},
    ]
    
    resp = client.chat.completions.create(
        model=model,
        messages=messages,
        tools=TOOLS,       # 上面定义的 Function Calling Schema
        temperature=0.3,
    )
    return resp

# 示例:Agent 会先调 get_realtime_ticker,拿到事实后再回答
result = ask_with_tools("688256.SH 现在多少钱?")

System Prompt 里的三句话是关键

  1. "必须先调用工具获取实时数据" —— 强制模型先取事实
  2. "不要用训练数据里的历史价格回答" —— 明确禁止幻觉来源
  3. "工具返回空数组,明确告知未找到" —— 处理不存在标的(实测返回 HTTP 200 + 空数组)

7.2 扣子工作流:把数据底座编排成 Agent 能力

扣子是字节的 Agent 开发平台,适合把多个工具、多个数据源编排成一个完整的 Agent 工作流。

一个典型的扣子工作流

用户提问:"688256.SH 最近走势怎么样,估值高不高?"
    ↓
[节点1] 意图识别:涉及实时价格 + 历史K线 + 估值
    ↓
[节点2] 并行工具调用:
    ├── get_realtime_ticker("688256.SH")
    ├── get_kline("688256.SH", 30天前, 今天)
    └── get_valuation("688256.SH")
    ↓
[节点3] 数据汇总:拼成结构化上下文
    ↓
[节点4] 豆包推理:基于事实生成分析
    ↓
[节点5] 输出:带数据来源的分析结论

扣子里配置 MCP 工具

  1. 进入「插件」→「创建插件」→「MCP 插件」
  2. 填入 MCP Server 地址(如果部署在 ECS 上,用公网或 VPC 内网地址)
  3. 测试工具调用是否正常
  4. 在工作流里引用这些工具

扣子的工作流设计要点

  • 每个工具节点必须有失败分支。工具调用失败时,Agent 要明确告知用户"数据获取失败",而不是编造。
  • 并行调用能省时间。实时快照、历史K线、估值三个工具没有依赖关系,并行调用能显著降低响应延迟。
  • 输出必须带数据来源和时间戳。让用户知道"这个数字是什么时候取的",避免误导。

7.3 先取事实,再推理

这是 AI 层的核心原则。模型只负责"基于给定事实推理",不负责"提供事实"。 事实永远来自工具调用。

# 完整流程示例
def agent_workflow(user_question: str, symbol: str):
    # 第一步:取事实
    facts = {}
    try:
        facts["ticker"] = get_ticker(symbol)
    except Exception as e:
        facts["ticker_error"] = str(e)
    
    try:
        now_ms = int(time.time() * 1000)
        month_ago_ms = now_ms - 30 * 24 * 3600 * 1000
        facts["klines"] = get_kline(symbol, month_ago_ms, now_ms)
    except Exception as e:
        facts["klines_error"] = str(e)
    
    # 第二步:把事实拼成上下文
    context = json.dumps(facts, ensure_ascii=False, indent=2)
    
    # 第三步:模型基于事实推理
    prompt = f"""以下是 {symbol} 的实时数据:
{context}

请基于以上事实回答:{user_question}
如果某项数据获取失败,明确告知,不要编造。"""
    
    return ask_doubao(prompt)

顺序反了,分析结论就没有可追溯的数据基础。


八、部署与监控:veFaaS + VKE + TLS

8.1 分层部署

组件部署方式原因
定时快照采集veFaaS无状态、按调用计费
WebSocket 常驻订阅ECS / 轻量服务器需要长连接
MCP ServerECS / VKE需要常驻进程
因子计算任务VKE CronJob需要弹性算力
Agent 编排火山方舟 / 扣子托管服务

8.2 VKE 因子 CronJob

apiVersion: batch/v1
kind: CronJob
metadata:
  name: factor-daily
spec:
  schedule: "0 16 * * 1-5"   # 工作日 16:00,A股收盘后
  jobTemplate:
    spec:
      template:
        spec:
          restartPolicy: OnFailure
          containers:
            - name: factor
              image: your-registry/factor:latest
              envFrom:
                - secretRef:
                    name: tickdb-secret
              resources:
                requests: { cpu: "500m", memory: "1Gi" }
                limits:   { cpu: "2",    memory: "4Gi" }

8.3 TLS 日志采集

veFaaS、ECS、VKE 的日志统一接到 TLS:

来源接入方式
veFaaS函数配置里开启日志投递,选择 TLS 日志项目
ECS装 TLS 日志采集器,采集 /var/log/tickdb/*.log
VKE用 TLS 的容器日志采集,按 Pod 标签过滤

统一日志格式:

{
  "level": "info",
  "stage": "mcp_tool_call",
  "tool": "get_realtime_ticker",
  "symbol": "688256.SH",
  "latency_ms": 42,
  "msg": "tool call success",
  "ts": "2026-09-16T09:30:00.123Z"
}

关键字段stage(在五层链路的哪一层)、tool(Agent 调用了哪个工具)、latency_ms(工具调用延迟)、symbol。四条齐全,Agent 出问题能定位到具体环节。

8.4 关键告警

告警条件处理
采集失败veFaaS 错误率 > 5% / 5 分钟检查 API Key 和网络
MCP 工具失败工具调用错误率 > 10% / 5 分钟检查 API 和参数
消息堆积Kafka lag > 10000检查消费者健康
Agent 幻觉检测工具未调用但回答含具体价格检查 System Prompt
数据延迟ingest_ts - data_ts > 60 秒检查采集端

最后一条"Agent 幻觉检测"是 Agent 场景特有的。如果 Agent 回答里出现了具体价格但日志里没有对应的工具调用记录,说明模型在编造。 这类告警能在 Prompt 设计失误时第一时间发现。


九、数据接入检查清单

链路项目已完成备注
接入REST 快照采集veFaaS 定时
接入WebSocket 常驻订阅ECS 长连接
接入MCP Server 封装工具层核心
接入密钥存环境变量禁止硬编码
传输Kafka Topic 规划分区键
传输幂等生产acks=all
存储TOS 冷数据归档Parquet + 生命周期
存储veDB 热数据分区 + 索引
工具MCP 工具描述清晰给模型看的提示
工具工具返回值结构化JSON 字符串
AgentSystem Prompt 强制先调工具禁止用记忆答价格
Agent空数组处理明确告知未找到
Agent输出带数据来源可追溯
监控TLS 日志统一采集含 tool 字段
监控Agent 幻觉告警无工具调用但有具体价格

十、总结

金融 Agent 的最大问题不是模型能力,是数据底座。 模型负责推理,数据负责事实。让模型用记忆猜价格,是把分析过程变成幻觉生成过程。

在火山引擎上搭这套底座,五层各司其职:

  • 接入层:veFaaS 定时拉快照,ECS 常驻订阅,密钥不落代码。
  • 传输层:Kafka 用 symbol 做分区键,acks=all 加幂等。
  • 存储层:TOS 归档冷数据,veDB 存热数据,冷热分离。
  • 工具层:MCP Server 把行情接口封装成 Agent 能调的工具,工具描述要写清使用场景。
  • Agent 层:火山方舟的豆包通过 Tool Call 取事实,扣子编排工作流,System Prompt 必须强制"先调工具再回答"

我在 Agent 数据底座上踩坑最深的三个点

  1. K线时间参数不是 beg_day/end_day,是 start_time/end_time 毫秒值。200 不等于参数生效。
  2. WebSocket 订阅用 cmd/data/symbols,不是 action/channel/symbol;close 1008 必须停止重连。
  3. Agent 不调工具就直接答价格,是 Prompt 设计问题,不是模型问题。 System Prompt 里必须明确禁止用训练数据回答实时问题。

你现在就能做的一件事:打开你的 Agent 对话记录,找一条涉及具体价格的回答,看日志里有没有对应的工具调用。如果没有,你的 Agent 正在猜价格。


文中 API 行为实测于 2026-09-16,样本标的为 688256.SH(寒武纪)和 600519.SH(贵州茅台)。火山引擎产品版本以实际部署环境为准。样本标的仅用于技术演示,不构成投资建议。

0
0
0
0
评论
未登录
暂无评论