实时新闻导火线系统
策略简介
本策略是一套运行在 FMZ 平台上的“新闻 + K 线 + 手动交易”看盘系统。它并不是自动判断利好利空并直接下单的量化策略,而是把实时新闻、行情 K 线、账户权益、持仓状态和交易按钮整合到同一个面板中,帮助用户更快发现新闻与价格波动之间的关系。
策略名称中的“导火线”,指的是很多剧烈行情并非单纯由技术指标推动,而是由突发新闻、政策表态、地缘冲突、宏观数据或重要人物发言引爆。系统的目标就是把这些可能引发行情的新闻直接标记到 K 线上,让用户在看价格变化时,也能同步看到“这一段行情可能因为什么发生”。
核心逻辑
-
接入实时新闻源
策略通过金十 MCP 接口获取实时快讯和资讯,包括
list_flash与list_news两类数据。新闻数据会被统一整理为标准格式:
字段 说明 ts新闻时间戳 time可读时间 title新闻标题 source新闻来源 full_text完整文本内容 这样后续图表、关键词过滤和状态面板不需要关心不同新闻源的原始字段差异。
-
将新闻标记到 K 线上
策略会在行情图表中绘制 K 线,并额外添加一层
flags新闻标记。每条命中条件的新闻会按照时间戳对齐到对应的 K 线 bar 上。用户把鼠标移动到标记上,就可以看到该时间点对应的新闻标题。
这样价格和新闻不再分散在两个窗口里,而是出现在同一张图上。
-
关键词过滤新闻
策略支持通过
News Keywords参数设置关键词,并用|分隔多个关键词,例如:text伊朗|加息|非农|关税|特朗普|美联储只有命中关键词的新闻才会优先标记到图表上,避免无关快讯过多导致图表被刷屏。
-
行情、账户、持仓一屏展示
除了 K 线和新闻标记,策略还会在状态面板中展示:
模块 内容 行情信息 当前交易品种、最新价格、K 线周期 账户信息 权益、余额、盈亏等 持仓信息 当前多空持仓情况 关键词新闻 命中关键词的相关新闻 最新快讯 最近获取到的新闻列表 用户可以在一个面板中同时观察行情、新闻和账户状态。
-
提供手动交易指令
策略保留“人在回路中”的交易方式,用户可以通过 FMZ 指令手动执行交易。
指令 作用 openLong开多 openShort开空 closeLong平多 closeShort平空 closeAll一键全平 amount修改默认下单数量 refreshNews手动刷新新闻 策略不替用户判断新闻方向,而是提供更集中的信息面板和更快捷的执行入口。
使用流程
text
接入金十 MCP 新闻源
↓
拉取实时快讯和资讯
↓
统一新闻格式并去重
↓
根据关键词筛选重要新闻
↓
将新闻按时间标记到 K 线上
↓
用户结合价格、新闻、持仓判断行情
↓
通过手动指令开仓、平仓或全平
主要参数
| 参数 | 说明 |
|---|---|
Symbol | 交易品种 |
Kline Period | K 线周期,默认 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、监管、利率、特朗普”等。
- 新闻刷新间隔不宜过长,否则会降低突发事件响应速度。
- 下单数量建议保持保守,尤其是在新闻驱动的高波动行情中。
- 可先作为纯看盘工具使用,熟悉新闻标记与行情反应后,再结合手动交易。
- 对特别重要的宏观事件,建议同时结合其他新闻源确认,避免单一来源误导。
相关文章
-
策略源码:实时新闻导火线系统
包含完整 FMZ Python 策略代码,可查看金十 MCP 接入、新闻归一化、K 线图表绘制、新闻 flags 标记、状态面板和手动交易指令实现。 -
策略文章:导火线:把“新闻”画进 K 线里,是一次什么样的尝试
文章详细介绍了为什么要把新闻与 K 线放在同一个视图中、如何通过 MCP 接入金十数据、如何将新闻按时间标记到 K 线上,以及该系统作为“信息整合 + 手动执行”工具的定位与局限。
# -*- 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)- 1