Type/to search

实时新闻导火线系统

Common strategy
Created: 2026-06-11 23:55:40
Last modified: a month ago
2
Follow
502
Followers

实时新闻导火线系统

策略简介

本策略是一套运行在 FMZ 平台上的“新闻 + K 线 + 手动交易”看盘系统。它并不是自动判断利好利空并直接下单的量化策略,而是把实时新闻、行情 K 线、账户权益、持仓状态和交易按钮整合到同一个面板中,帮助用户更快发现新闻与价格波动之间的关系。

策略名称中的“导火线”,指的是很多剧烈行情并非单纯由技术指标推动,而是由突发新闻、政策表态、地缘冲突、宏观数据或重要人物发言引爆。系统的目标就是把这些可能引发行情的新闻直接标记到 K 线上,让用户在看价格变化时,也能同步看到“这一段行情可能因为什么发生”。

核心逻辑

  1. 接入实时新闻源

    策略通过金十 MCP 接口获取实时快讯和资讯,包括 list_flashlist_news 两类数据。

    新闻数据会被统一整理为标准格式:

    字段说明
    ts新闻时间戳
    time可读时间
    title新闻标题
    source新闻来源
    full_text完整文本内容

    这样后续图表、关键词过滤和状态面板不需要关心不同新闻源的原始字段差异。

  2. 将新闻标记到 K 线上

    策略会在行情图表中绘制 K 线,并额外添加一层 flags 新闻标记。

    每条命中条件的新闻会按照时间戳对齐到对应的 K 线 bar 上。用户把鼠标移动到标记上,就可以看到该时间点对应的新闻标题。

    这样价格和新闻不再分散在两个窗口里,而是出现在同一张图上。

  3. 关键词过滤新闻

    策略支持通过 News Keywords 参数设置关键词,并用 | 分隔多个关键词,例如:

    text
    伊朗|加息|非农|关税|特朗普|美联储

    只有命中关键词的新闻才会优先标记到图表上,避免无关快讯过多导致图表被刷屏。

  4. 行情、账户、持仓一屏展示

    除了 K 线和新闻标记,策略还会在状态面板中展示:

    模块内容
    行情信息当前交易品种、最新价格、K 线周期
    账户信息权益、余额、盈亏等
    持仓信息当前多空持仓情况
    关键词新闻命中关键词的相关新闻
    最新快讯最近获取到的新闻列表

    用户可以在一个面板中同时观察行情、新闻和账户状态。

  5. 提供手动交易指令

    策略保留“人在回路中”的交易方式,用户可以通过 FMZ 指令手动执行交易。

    指令作用
    openLong开多
    openShort开空
    closeLong平多
    closeShort平空
    closeAll一键全平
    amount修改默认下单数量
    refreshNews手动刷新新闻

    策略不替用户判断新闻方向,而是提供更集中的信息面板和更快捷的执行入口。

使用流程

text
接入金十 MCP 新闻源 ↓ 拉取实时快讯和资讯 ↓ 统一新闻格式并去重 ↓ 根据关键词筛选重要新闻 ↓ 将新闻按时间标记到 K 线上 ↓ 用户结合价格、新闻、持仓判断行情 ↓ 通过手动指令开仓、平仓或全平

主要参数

参数说明
Symbol交易品种
Kline PeriodK 线周期,默认 1 分钟
Kline Limit图表加载的 K 线数量
Refresh Interval(s)行情和面板刷新间隔
Default Order Amount默认下单数量
Jin10 MCP Token金十 MCP Token,可选但接入新闻源时需要
Jin10 MCP URL金十 MCP 接口地址
News Keywords新闻关键词,多个关键词用 | 分隔
News Refresh Interval(s)新闻刷新间隔

策略特点

  • 直接运行在 FMZ 平台,无需额外部署外部服务。
  • 将实时新闻标记到 K 线图上,方便观察新闻与价格波动的时间关系。
  • 支持金十 MCP 快讯和资讯接口。
  • 支持关键词过滤,减少无关信息干扰。
  • 同屏展示行情、账户、持仓和新闻。
  • 提供开多、开空、平仓、一键全平等手动交易指令。
  • 决策权保留给用户,适合辅助看盘和事件驱动交易。

适用场景

本策略适合关注突发新闻、宏观数据、地缘冲突、政策表态、央行讲话等事件驱动行情的用户。

例如原油、黄金、BTC、股指、外汇等品种,经常会受到新闻快速影响。单看 K 线只能看到“价格发生了什么”,而新闻标记可以帮助用户回到当时的信息现场,理解“为什么这一根 K 线突然拉升或跳水”。

它更像是一个交易辅助面板,而不是全自动交易策略。适合有一定主观判断能力、希望提高看盘效率和事件响应速度的用户。

风险提示

  • 策略不会自动判断新闻是利好还是利空,最终交易决策仍由用户完成。
  • 新闻与 K 线的对应关系只是按时间戳粗略对齐,不代表价格波动一定由该新闻导致。
  • 关键词命中不等于新闻重要,未命中也不代表新闻无效。
  • MCP 接口可能存在延迟、请求失败、Token 失效或数据缺失。
  • 新闻行情反应速度可能非常快,人工下单可能错过最佳价格。
  • 手动交易仍存在追涨杀跌、误判消息、滑点和流动性风险。
  • 一键全平等指令需要谨慎使用,避免误操作。

使用建议

  • 根据交易品种设置关键词,例如原油关注“伊朗、OPEC、库存、制裁”,黄金关注“美联储、降息、CPI、非农”,BTC 关注“ETF、监管、利率、特朗普”等。
  • 新闻刷新间隔不宜过长,否则会降低突发事件响应速度。
  • 下单数量建议保持保守,尤其是在新闻驱动的高波动行情中。
  • 可先作为纯看盘工具使用,熟悉新闻标记与行情反应后,再结合手动交易。
  • 对特别重要的宏观事件,建议同时结合其他新闻源确认,避免单一来源误导。

相关文章

Source
Python
# -*- coding: utf-8 -*-
"""
FUSE — Manual News Trading Panel (FMZ Python)
直接在 FMZ 平台运行,无需外部服务。
"""

import json
import time
import urllib.error
import urllib.request
import calendar

# ─── 常量映射 ─────────────────────────────────────────────────────────────────

PERIOD_MAP = {
    "1m":  PERIOD_M1,  "3m":  PERIOD_M3,  "5m":  PERIOD_M5,
    "15m": PERIOD_M15, "30m": PERIOD_M30,
    "1h":  PERIOD_H1,  "4h":  PERIOD_H4,  "1d":  PERIOD_D1,
}

PERIOD_MS = {
    "1m":  60000,    "3m":  180000,   "5m":  300000,
    "15m": 900000,   "30m": 1800000,
    "1h":  3600000,  "4h":  14400000, "1d":  86400000,
}

# ─── 状态 ─────────────────────────────────────────────────────────────────────

_chart          = None
_last_symbol    = ""
_cur_amount     = ORDER_AMOUNT
_cached_news    = []
_last_news_at   = 0
_mcp_session_id = ""
_mcp_req_id     = 0
_mcp_ready      = False
_init_equity    = None


# ─── MCP Client ───────────────────────────────────────────────────────────────

def _mcp_headers():
    h = {
        "Content-Type": "application/json",
        "Accept":        "application/json, text/event-stream",
        "Authorization": "Bearer " + JIN10_MCP_TOKEN,
        "User-Agent":    "FUSE-FMZ/1.0",
    }
    if _mcp_session_id:
        h["Mcp-Session-Id"] = _mcp_session_id
    return h


def _mcp_parse(text):
    last = None
    for line in text.splitlines():
        line = line.strip()
        if not line.startswith("data:"):
            continue
        chunk = line[5:].strip()
        if not chunk or chunk == "[DONE]":
            continue
        try:
            parsed = json.loads(chunk)
            if isinstance(parsed, dict):
                last = parsed
        except json.JSONDecodeError:
            pass
    if last is None:
        try:
            last = json.loads(text)
        except json.JSONDecodeError:
            return {"raw": text[:2000]}
    if "error" in last:
        raise RuntimeError("MCP error: " + str(last["error"]))
    return last.get("result", last)


def _mcp_post(payload, is_notification=False):
    global _mcp_session_id, _mcp_req_id
    if not is_notification:
        _mcp_req_id += 1
        payload["id"] = _mcp_req_id
    body = json.dumps(payload, ensure_ascii=False).encode()
    req  = urllib.request.Request(
        JIN10_MCP_URL, data=body, headers=_mcp_headers(), method="POST"
    )
    try:
        with urllib.request.urlopen(req, timeout=30) as resp:
            sid = resp.headers.get("Mcp-Session-Id")
            if sid:
                _mcp_session_id = sid
            if resp.status == 202:
                return {}
            text = resp.read().decode("utf-8", errors="replace")
    except urllib.error.HTTPError as e:
        raise RuntimeError("HTTP %d: %s" % (e.code, e.read().decode()[:400]))
    except urllib.error.URLError as e:
        raise RuntimeError("Network: " + str(e))
    return _mcp_parse(text)


def mcp_rpc(method, params=None):
    return _mcp_post({"jsonrpc": "2.0", "method": method, "params": params or {}})


def mcp_notify(method):
    _mcp_post({"jsonrpc": "2.0", "method": method, "params": {}}, is_notification=True)


def mcp_init():
    global _mcp_ready
    mcp_rpc("initialize", {
        "protocolVersion": "2025-11-25",
        "capabilities":    {},
        "clientInfo":      {"name": "fuse-fmz", "version": "1.0"},
    })
    mcp_notify("notifications/initialized")
    _mcp_ready = True
    Log("MCP ready  session_id=" + (_mcp_session_id or "(none)"))


def mcp_call_tool(name, arguments=None):
    result = mcp_rpc("tools/call", {"name": name, "arguments": arguments or {}})
    if "structuredContent" in result:
        return result["structuredContent"]
    for item in result.get("content", []):
        if item.get("type") != "text":
            continue
        try:
            return json.loads(item["text"])
        except json.JSONDecodeError:
            return {"raw": item["text"]}
    return result

# ─── News helpers ─────────────────────────────────────────────────────────────

def _extract_items(payload):
    if isinstance(payload, list):
        return [x for x in payload if isinstance(x, dict)]
    if not isinstance(payload, dict):
        return []
    for key in ("items", "data", "list", "rows", "flash", "news", "result"):
        v = payload.get(key)
        if isinstance(v, list):
            return [x for x in v if isinstance(x, dict)]
        if isinstance(v, dict):
            nested = _extract_items(v)
            if nested:
                return nested
    return []


def _pick(obj, *keys):
    for k in keys:
        v = obj.get(k)
        if v is not None and str(v).strip():
            return str(v).strip()
    return ""


import re

def _to_ms(ts_str):
    if not ts_str:
        return 0
    m = re.match(r"(\d{4}-\d{2}-\d{2})T(\d{2}:\d{2}:\d{2})([+-])(\d{2}):?(\d{2})?", ts_str.strip())
    if m:
        date_part, time_part, sign, tzh, tzm = m.groups()
        t = time.strptime(date_part + " " + time_part, "%Y-%m-%d %H:%M:%S")
        ts = calendar.timegm(t)
        offset = int(tzh) * 3600 + int(tzm or 0) * 60
        ts = ts - offset if sign == "+" else ts + offset
        return ts * 1000
    try:
        n = float(ts_str)
        return int(n * 1000 if n < 1e12 else n)
    except Exception:
        return 0


def _normalize(items, source):
    rows = []
    for it in items:
        title = _pick(it, "title", "content", "summary", "introduction", "name")
        if not title:
            continue
        ts_str = _pick(it, "time", "ts", "time_str", "datetime", "created_at", "display_time")
        full_text = " ".join(str(it.get(f, "")) for f in ("title", "content", "introduction", "summary"))
        rows.append({"ts": _to_ms(ts_str), "time": ts_str, "title": title, "source": source, "full_text": full_text})
    rows.sort(key=lambda x: x["ts"], reverse=True)
    return rows


def refresh_news():
    global _cached_news, _last_news_at, _mcp_ready
    if not JIN10_MCP_TOKEN:
        return
    now = int(time.time())
    if now - _last_news_at < NEWS_REFRESH_SEC and _cached_news:
        return
    _last_news_at = now
    try:
        if not _mcp_ready:
            mcp_init()
        flash_raw = mcp_call_tool("list_flash")
        news_raw  = mcp_call_tool("list_news")
        combined  = (
            _normalize(_extract_items(flash_raw), "flash") +
            _normalize(_extract_items(news_raw),  "news")
        )
        combined.sort(key=lambda x: x["ts"], reverse=True)
        _cached_news = combined[:80]
        Log("News updated: %d items" % len(_cached_news))
    except Exception as e:
        Log("News refresh failed: " + str(e))
        _mcp_ready = False

# ─── Symbol helpers ───────────────────────────────────────────────────────────

def switch_symbol(symbol):
    global _last_symbol
    if _last_symbol == symbol:
        return
    parts = symbol.split(".")
    exchange.SetCurrency(parts[0])
    exchange.SetContractType(parts[1] if len(parts) > 1 else "swap")
    _last_symbol = symbol


def get_records(symbol):
    switch_symbol(symbol)
    recs = exchange.GetRecords(symbol, PERIOD_MAP.get(KLINE_PERIOD, PERIOD_M1), KLINE_LIMIT)
    if not recs:
        return []
    recs = list(recs[-KLINE_LIMIT:])
    for r in recs:
        if r['Time'] < 1e12:
            r['Time'] = int(r['Time'] * 1000)
    return recs


def get_ticker(symbol):
    switch_symbol(symbol)
    return exchange.GetTicker(symbol)


def get_positions(symbol):
    switch_symbol(symbol)
    pos = exchange.GetPositions(symbol)
    return list(pos) if pos else []


def get_equity():
    try:
        acc = exchange.GetAccount()
        if not acc:
            return None
        return float(acc.Equity)
    except Exception as e:
        Log("get_equity error: " + str(e))
        return None

# ─── Equity tracking ──────────────────────────────────────────────────────────

def init_equity():
    global _init_equity
    saved = _G("init_equity")
    if saved is not None:
        _init_equity = float(saved)
        Log("Restored init_equity: %.4f" % _init_equity)
        return
    eq = get_equity()
    if eq is not None:
        _init_equity = eq
        _G("init_equity", eq)
        Log("Init equity recorded: %.4f" % eq)


def record_profit(equity):
    if _init_equity is None or equity is None:
        return
    LogProfit(round(equity - _init_equity, 4), '&')

# ─── Trade helpers ────────────────────────────────────────────────────────────

def market_order(symbol, action, amount):
    switch_symbol(symbol)
    qty = float(amount or _cur_amount)
    dir_map = {
        "openLong":   ("buy",       "open long"),
        "openShort":  ("sell",      "open short"),
        "closeLong":  ("closebuy",  "close long"),
        "closeShort": ("closesell", "close short"),
    }
    if action not in dir_map:
        return
    direction, memo = dir_map[action]
    Log("Trade: %s %s %s" % (action, qty, symbol))
    exchange.CreateOrder(symbol, direction, -1, qty, memo)


def close_all(symbol):
    for p in get_positions(symbol):
        qty = abs(float(p.Amount))
        if qty <= 0:
            continue
        direction = "closesell" if p.Type == PD_SHORT else "closebuy"
        Log("CloseAll: %s %s" % (direction, qty))
        exchange.CreateOrder(symbol, direction, -1, qty, "close all")

# ─── Command handler ──────────────────────────────────────────────────────────

def handle_command(symbol):
    global _cur_amount, _last_news_at
    cmd = GetCommand()
    if not cmd:
        return
    Log("CMD: " + cmd)
    parts = cmd.split(":")
    key   = parts[0]
    val   = parts[1] if len(parts) > 1 else ""
    if   key == "openLong":    market_order(symbol, "openLong",   _cur_amount)
    elif key == "openShort":   market_order(symbol, "openShort",  _cur_amount)
    elif key == "closeLong":   market_order(symbol, "closeLong",  _cur_amount)
    elif key == "closeShort":  market_order(symbol, "closeShort", _cur_amount)
    elif key == "closeAll":    close_all(symbol)
    elif key == "amount":
        _cur_amount = float(val)
        Log("Amount updated: " + str(_cur_amount))

# ─── Chart ────────────────────────────────────────────────────────────────────

def init_chart(symbol):
    global _chart
    _chart = Chart({
        "__isStock": True,
        "chart":     {"style": {"fontFamily": "Microsoft YaHei, SimHei, Arial, sans-serif"}},
        "title":     {"text": "FUSE  " + symbol},
        "xAxis":     {"type": "datetime"},
        "series": [
            {
                "id":   "kline",
                "type": "candlestick",
                "name": symbol,
                "data": [],
            },
            {
                "type":      "flags",
                "name":      "News",
                "onSeries":  "kline",
                "shape":     "circlepin",
                "color":     "#F59E0B",
                "fillColor": "#F59E0B",
                "width":     16,
                "data":      [],
            },
        ],
    })
    _chart.reset()


_last_bar_time   = 0
_last_news_hash  = 0   # 新闻变化时触发图表重置
_flagged_news_ts = set()  # 已绘制过的新闻 ts,避免重复 add

def draw_chart(records):
    global _last_bar_time, _last_news_hash, _flagged_news_ts

    if not _chart or not records:
        return

    # 检测新闻是否有更新,有则重置图表重画
    news_hash = hash(tuple(n.get("ts", 0) for n in _cached_news[:10]))
    news_changed = (news_hash != _last_news_hash)
    if news_changed:
        _chart.reset()
        _last_bar_time   = 0
        _last_news_hash  = news_hash
        _flagged_news_ts = set()

    # series 0:K线,增量 add
    for r in records:
        t   = r['Time']
        bar = [t, r['Open'], r['High'], r['Low'], r['Close']]
        if t > _last_bar_time:
            _chart.add(0, bar)
            _last_bar_time = t
        elif t == _last_bar_time:
            _chart.add(0, bar, -1)

    # series 1:关键词新闻 flag,对应到K线 bar 时间
    if not _cached_news:
        return
    kws = [k.strip() for k in NEWS_KEYWORD.split("|") if k.strip()]
    kw_news = [n for n in _cached_news if not kws or any(k in n.get("full_text", n["title"]) for k in kws)]
    if not kw_news:
        return
    p_ms  = PERIOD_MS.get(KLINE_PERIOD, 60000)
    first = records[0]['Time']
    last  = records[-1]['Time']
    by_bar = {}
    for n in kw_news:
        if not n.get("ts"):
            continue
        key = (n["ts"] // p_ms) * p_ms
        if key not in by_bar:
            by_bar[key] = n

    for ts, item in sorted(by_bar.items()):
        if not (first <= ts <= last):
            continue
        if ts in _flagged_news_ts:
            continue
        _chart.add(1, {
            "x": ts,
            "title": "📰",
            "text": item["title"][:100],
        })
        _flagged_news_ts.add(ts)

# ─── Status panel ─────────────────────────────────────────────────────────────

def make_status(symbol, ticker, positions, equity):
    # 1. 行情
    quote_rows = [[
        symbol,
        ticker.Last if ticker else "-",
        ticker.High if ticker else "-",
        ticker.Low  if ticker else "-",
        _D(ticker.Time) if ticker and ticker.Time else _D(),
    ]]

    # 2. 账户权益
    if equity is not None and _init_equity is not None:
        pnl = round(equity - _init_equity, 4)
        roi = round(pnl / _init_equity * 100, 2) if _init_equity else 0
        acc_rows = [[round(_init_equity, 4), round(equity, 4), pnl, str(roi) + "%"]]
    else:
        acc_rows = [["-", "-", "-", "-"]]

    # 3. 持仓
    pos_rows = []
    if positions:
        for p in positions:
            side = "空 SHORT" if p.Type == PD_SHORT else "多 LONG"
            pnl  = "-"
            if ticker and ticker.Last and p.Price:
                raw = (
                    (p.Price - ticker.Last) * abs(p.Amount)
                    if p.Type == PD_SHORT
                    else (ticker.Last - p.Price) * abs(p.Amount)
                )
                pnl = round(raw, 4)
            pos_rows.append([side, round(p.Amount, 4), round(p.Price, 4), pnl])
    else:
        pos_rows.append(["无持仓 / Flat", "-", "-", "-"])

    # 4. 关键词新闻(显示命中的关键词,而不是来源)
    kws = [k.strip() for k in NEWS_KEYWORD.split("|") if k.strip()]
    kw_rows = []
    for item in _cached_news[:40]:
        t     = item.get("time") or (_D(item["ts"]) if item.get("ts") else "-")
        title = item["title"][:90]
        text  = item.get("full_text", item["title"])
        hit_kws = [k for k in kws if k in text]
        if hit_kws:
            kw_rows.append([t, "/".join(hit_kws), title])

    if not kw_rows:
        kw_rows = [["-", "-", "暂无关键词相关新闻"]]

    # 5. 实时新闻(最新5条)
    all_rows = []
    for item in _cached_news[:5]:
        t     = item.get("time") or (_D(item["ts"]) if item.get("ts") else "-")
        src   = item.get("source", "-")
        title = item["title"][:90]
        all_rows.append([t, src, title])
    if not all_rows:
        all_rows = [["-", "-", "新闻加载中..." if JIN10_MCP_TOKEN else "未配置 MCP Token"]]

    LogStatus(
        "`" + json.dumps({
            "type": "table", "title": "行情 Quotes",
            "cols": ["Symbol", "Last", "High", "Low", "Time"],
            "rows": quote_rows,
        }) + "`\n" +
        "`" + json.dumps({
            "type": "table", "title": "账户权益 Equity",
            "cols": ["初始权益 Init", "当前权益 Now", "盈亏 PnL", "收益率 ROI"],
            "rows": acc_rows,
        }) + "`\n" +
        "`" + json.dumps({
            "type": "table", "title": "持仓 Positions",
            "cols": ["方向 Side", "数量 Amount", "均价 AvgPrice", "浮盈 UPnL"],
            "rows": pos_rows,
        }) + "`\n" +
        "`" + json.dumps({
            "type": "table", "title": "关键词新闻 (" + NEWS_KEYWORD + ")",
            "cols": ["时间", "关键词", "标题"],
            "rows": kw_rows,
        }) + "`\n" +
        "`" + json.dumps({
            "type": "table", "title": "实时新闻 Live News (最新5条)",
            "cols": ["时间", "来源", "标题"],
            "rows": all_rows,
        }) + "`"
    )

# ─── Main ─────────────────────────────────────────────────────────────────────

def main():
    _G(None)
    if not SYMBOL:
        raise RuntimeError("SYMBOL is empty")

    # 先切换合约,确保 GetAccount 返回正确权益
    switch_symbol(SYMBOL)

    # 恢复或记录初始权益
    init_equity()

    # 初始化 MCP session
    if JIN10_MCP_TOKEN:
        try:
            mcp_init()
        except Exception as e:
            Log("MCP init failed: " + str(e))

    init_chart(SYMBOL)
    Log("FUSE started | %s | %s | amount=%s" % (SYMBOL, KLINE_PERIOD, _cur_amount))

    while True:
        handle_command(SYMBOL)
        refresh_news()

        records   = get_records(SYMBOL)
        ticker    = get_ticker(SYMBOL)
        positions = get_positions(SYMBOL)
        equity    = get_equity()

        record_profit(equity)

        if records:
            draw_chart(records)

        make_status(SYMBOL, ticker, positions, equity)

        Sleep(max(1, REFRESH_SECONDS) * 1000)
Strategy parameters
Strategy parameters
Symbol
Kline Period
Kline Limit
Refresh Interval(s)
Default Order Amount
Jin10 MCP Token (Optional)
Jin10 MCP URL (Optional)
News Keywords (Optional)
News Refresh Interval(s)
Commands
Open Long
openLong
Open Short
openShort
Close Long
closeLong
Close Short
closeShort
Close All
closeAll
Set Amount
amount
Refresh News
refreshNews
Comment
All comments (0)
No data
No data
  • 1
Forums
PINE Language
Get the app
iPhone Download
© 2015 - ∞ INVENTOR PTE LTD (SG)