From f6bdc5a983af132dd4c3e959219c51f036b1e579 Mon Sep 17 00:00:00 2001 From: lookt Date: Tue, 15 Sep 2026 20:09:46 +0800 Subject: [PATCH] =?UTF-8?q?linux-web:=20B/S=20=E6=9E=B6=E6=9E=84=E2=80=94?= =?UTF-8?q?=E2=80=94=E9=87=87=E9=9B=86=E5=BC=95=E6=93=8E=E5=90=8E=E5=8F=B0?= =?UTF-8?q?=E7=BA=BF=E7=A8=8B=20+=20aiohttp=20WebSocket=20=E5=AE=9E?= =?UTF-8?q?=E6=97=B6=E6=8E=A8=E9=80=81=20+=20Web=20=E4=BB=AA=E8=A1=A8?= =?UTF-8?q?=E7=9B=98=20+=20systemd=20=E9=83=A8=E7=BD=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- deploy/a-stock-timeline.service | 18 +++ docs/DEPLOY_LINUX.md | 51 ++++++ docs/overnight_transmission.md | 51 ++++++ requirements.txt | 8 + src/realtime/collector.py | 267 ++++++++++++++++++++++++++++++++ src/web/index.html | 125 +++++++++++++++ src/web/server.py | 110 +++++++++++++ start_server.sh | 3 + 8 files changed, 633 insertions(+) create mode 100644 deploy/a-stock-timeline.service create mode 100644 docs/DEPLOY_LINUX.md create mode 100644 docs/overnight_transmission.md create mode 100644 requirements.txt create mode 100644 src/realtime/collector.py create mode 100644 src/web/index.html create mode 100644 src/web/server.py create mode 100644 start_server.sh diff --git a/deploy/a-stock-timeline.service b/deploy/a-stock-timeline.service new file mode 100644 index 0000000..b50f2d8 --- /dev/null +++ b/deploy/a-stock-timeline.service @@ -0,0 +1,18 @@ +[Unit] +Description=JQuant A股时间线实时采集与分析服务(B/S) +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +# 部署路径按实际调整(git clone linux-web 分支后的目录) +WorkingDirectory=/opt/a_stock_timeline +ExecStart=/opt/a_stock_timeline/venv/bin/python -m src.web.server +Restart=always +RestartSec=5 +User=www-data +Environment=PORT=8100 +# 日志进 journald:journalctl -u a-stock-timeline -f + +[Install] +WantedBy=multi-user.target diff --git a/docs/DEPLOY_LINUX.md b/docs/DEPLOY_LINUX.md new file mode 100644 index 0000000..b6873f9 --- /dev/null +++ b/docs/DEPLOY_LINUX.md @@ -0,0 +1,51 @@ +# Linux B/S 部署(linux-web 分支) + +架构:Linux 服务器常驻执行「实时抓取 + 分析」(采集引擎每 60s 一轮), +分析事件经 **WebSocket** 实时推送到所有已连接的 Web 仪表盘(浏览器打开即看,无需刷新)。 + +``` +采集引擎 TimelineCollector(后台线程) + ├─ 盘前:美股隔夜收盘 → 实证先验(skill)预估今日跳空/日内概率 + ├─ 盘中:东财7x24快讯入库 → 关键词告警 → 上证实时点位 + └─ 盘后:龙虎榜入库 → 收盘日报 + │ 事件回调(每条:新快讯/盘中点位/告警/日报…) + ▼ +aiohttp 服务(端口 8100) + ├─ GET / Web 仪表盘(深色,事件流实时上屏) + ├─ GET /api/recent 最近事件 JSON + ├─ GET /api/stats 连接数/缓冲统计 + └─ WS /ws 实时广播(断线自动重连) +``` + +## 部署步骤 + +```bash +sudo mkdir -p /opt && cd /opt +sudo git clone -b linux-web <仓库地址> a_stock_timeline +cd a_stock_timeline + +python3 -m venv venv +venv/bin/pip install -r requirements.txt + +sudo cp deploy/a-stock-timeline.service /etc/systemd/system/ +sudo systemctl daemon-reload +sudo systemctl enable --now a-stock-timeline + +# 查看实时日志 +journalctl -u a-stock-timeline -f +``` + +浏览器访问 `http://<服务器IP>:8100`(生产建议前置 nginx 做 TLS/域名)。 + +## 数据说明 + +- 数据库:`/opt/a_stock_timeline/data/a_stock.db`(SQLite,与 C 端同 schema) +- 事件流水:`data/realtime_feed.jsonl`(重启后自动回放最近 200 条到新连接的客户端) +- 实证先验:`docs/overnight_transmission.md`(由 windows-desktop 分支的 + `src/analysis/mine_patterns.py` 挖掘生成,拷贝到 docs/ 即可被服务端加载) + +## 与 C 端(windows-desktop 分支)的关系 + +- 共享:fetcher/storage/分析引擎与全部实证口径 +- 差异:C 端入口为 `start_realtime.bat` + 控制台输出 + Windows 数据路径; + linux-web 端入口为 `python -m src.web.server`,事件走 WebSocket 广播,路径全部相对化 diff --git a/docs/overnight_transmission.md b/docs/overnight_transmission.md new file mode 100644 index 0000000..c967ed0 --- /dev/null +++ b/docs/overnight_transmission.md @@ -0,0 +1,51 @@ +## 隔夜传导矩阵(美股 -> A股) + +> 口径:美股T日涨跌幅分箱 -> A股T+1交易日(跳空=开盘/前收-1;日内=收盘/开盘-1)。样本 2005-2026(受 A 股指数历史与纳斯达克 2014 起数据限制)。 + +### 纳斯达克100T日 -> 上证指数T+1日 + +- 美股跌>1.5%(305天): A股跳空 N=305, 均值 -0.64%, 中位数 -0.41%, 胜率 10.5%; 日内 N=305, 均值 0.20%, 中位数 0.15%, 胜率 56.4%; 全天 N=305, 均值 -0.44%, 中位数 -0.25%, 胜率 37.0% +- 美股跌0.5~1.5%(484天): A股跳空 N=484, 均值 -0.20%, 中位数 -0.18%, 胜率 24.6%; 日内 N=484, 均值 0.06%, 中位数 0.10%, 胜率 56.4%; 全天 N=484, 均值 -0.14%, 中位数 -0.06%, 胜率 44.6% +- 美股正负0.5%内(1287天): A股跳空 N=1287, 均值 -0.07%, 中位数 -0.06%, 胜率 36.3%; 日内 N=1287, 均值 0.15%, 中位数 0.14%, 胜率 58.7%; 全天 N=1287, 均值 0.08%, 中位数 0.07%, 胜率 53.9% +- 美股涨0.5~1.5%(737天): A股跳空 N=737, 均值 0.07%, 中位数 0.04%, 胜率 56.3%; 日内 N=737, 均值 0.08%, 中位数 0.06%, 胜率 52.8%; 全天 N=737, 均值 0.15%, 中位数 0.09%, 胜率 56.3% +- 美股涨>1.5%(309天): A股跳空 N=309, 均值 0.22%, 中位数 0.17%, 胜率 73.8%; 日内 N=309, 均值 0.02%, 中位数 0.05%, 胜率 53.4%; 全天 N=309, 均值 0.24%, 中位数 0.23%, 胜率 61.5% + +### 纳斯达克100T日 -> 沪深300T+1日 + +- 美股跌>1.5%(305天): A股跳空 N=305, 均值 -0.69%, 中位数 -0.50%, 胜率 13.8%; 日内 N=305, 均值 0.23%, 中位数 0.13%, 胜率 56.4%; 全天 N=305, 均值 -0.47%, 中位数 -0.34%, 胜率 35.1% +- 美股跌0.5~1.5%(484天): A股跳空 N=484, 均值 -0.20%, 中位数 -0.20%, 胜率 24.2%; 日内 N=484, 均值 0.04%, 中位数 0.04%, 胜率 51.7%; 全天 N=484, 均值 -0.16%, 中位数 -0.14%, 胜率 41.5% +- 美股正负0.5%内(1287天): A股跳空 N=1287, 均值 -0.04%, 中位数 -0.05%, 胜率 42.0%; 日内 N=1287, 均值 0.13%, 中位数 0.10%, 胜率 55.7%; 全天 N=1287, 均值 0.09%, 中位数 0.05%, 胜率 52.5% +- 美股涨0.5~1.5%(737天): A股跳空 N=737, 均值 0.13%, 中位数 0.07%, 胜率 63.0%; 日内 N=737, 均值 0.04%, 中位数 0.01%, 胜率 50.5%; 全天 N=737, 均值 0.17%, 中位数 0.10%, 胜率 55.0% +- 美股涨>1.5%(309天): A股跳空 N=309, 均值 0.32%, 中位数 0.25%, 胜率 79.0%; 日内 N=309, 均值 -0.06%, 中位数 -0.08%, 胜率 45.6%; 全天 N=309, 均值 0.26%, 中位数 0.21%, 胜率 61.2% + +### 标普500T日 -> 上证指数T+1日 + +- 美股跌>1.5%(366天): A股跳空 N=366, 均值 -0.92%, 中位数 -0.72%, 胜率 7.7%; 日内 N=366, 均值 0.27%, 中位数 0.28%, 胜率 57.9%; 全天 N=366, 均值 -0.65%, 中位数 -0.40%, 胜率 36.1% +- 美股跌0.5~1.5%(868天): A股跳空 N=868, 均值 -0.23%, 中位数 -0.21%, 胜率 23.6%; 日内 N=868, 均值 0.02%, 中位数 0.08%, 胜率 53.8%; 全天 N=868, 均值 -0.22%, 中位数 -0.16%, 胜率 43.2% +- 美股正负0.5%内(2804天): A股跳空 N=2804, 均值 -0.06%, 中位数 -0.05%, 胜率 39.2%; 日内 N=2804, 均值 0.13%, 中位数 0.13%, 胜率 56.0%; 全天 N=2804, 均值 0.07%, 中位数 0.06%, 胜率 53.2% +- 美股涨0.5~1.5%(1218天): A股跳空 N=1218, 均值 0.08%, 中位数 0.06%, 胜率 61.7%; 日内 N=1218, 均值 0.10%, 中位数 0.10%, 胜率 55.1%; 全天 N=1218, 均值 0.19%, 中位数 0.15%, 胜率 59.2% +- 美股涨>1.5%(319天): A股跳空 N=319, 均值 0.53%, 中位数 0.37%, 胜率 83.4%; 日内 N=319, 均值 -0.09%, 中位数 -0.09%, 胜率 46.1%; 全天 N=319, 均值 0.44%, 中位数 0.28%, 胜率 62.1% + +### 标普500T日 -> 沪深300T+1日 + +- 美股跌>1.5%(366天): A股跳空 N=366, 均值 -0.98%, 中位数 -0.77%, 胜率 9.0%; 日内 N=366, 均值 0.35%, 中位数 0.21%, 胜率 57.1%; 全天 N=366, 均值 -0.64%, 中位数 -0.49%, 胜率 38.0% +- 美股跌0.5~1.5%(868天): A股跳空 N=868, 均值 -0.25%, 中位数 -0.25%, 胜率 24.9%; 日内 N=868, 均值 0.03%, 中位数 0.00%, 胜率 49.0%; 全天 N=868, 均值 -0.22%, 中位数 -0.20%, 胜率 42.1% +- 美股正负0.5%内(2804天): A股跳空 N=2804, 均值 -0.04%, 中位数 -0.04%, 胜率 43.8%; 日内 N=2804, 均值 0.13%, 中位数 0.04%, 胜率 51.9%; 全天 N=2804, 均值 0.08%, 中位数 0.06%, 胜率 52.6% +- 美股涨0.5~1.5%(1218天): A股跳空 N=1218, 均值 0.13%, 中位数 0.11%, 胜率 65.4%; 日内 N=1218, 均值 0.08%, 中位数 0.02%, 胜率 50.5%; 全天 N=1218, 均值 0.21%, 中位数 0.15%, 胜率 57.1% +- 美股涨>1.5%(319天): A股跳空 N=319, 均值 0.61%, 中位数 0.44%, 胜率 85.9%; 日内 N=319, 均值 -0.18%, 中位数 -0.17%, 胜率 40.8%; 全天 N=319, 均值 0.43%, 中位数 0.32%, 胜率 62.7% + +### 道琼斯T日 -> 上证指数T+1日 + +- 美股跌>1.5%(324天): A股跳空 N=324, 均值 -0.95%, 中位数 -0.78%, 胜率 8.6%; 日内 N=324, 均值 0.32%, 中位数 0.29%, 胜率 59.0%; 全天 N=324, 均值 -0.64%, 中位数 -0.43%, 胜率 37.0% +- 美股跌0.5~1.5%(890天): A股跳空 N=890, 均值 -0.24%, 中位数 -0.20%, 胜率 22.6%; 日内 N=890, 均值 0.11%, 中位数 0.13%, 胜率 56.4%; 全天 N=890, 均值 -0.14%, 中位数 -0.06%, 胜率 45.1% +- 美股正负0.5%内(2894天): A股跳空 N=2894, 均值 -0.06%, 中位数 -0.05%, 胜率 39.9%; 日内 N=2894, 均值 0.09%, 中位数 0.10%, 胜率 54.8%; 全天 N=2894, 均值 0.03%, 中位数 0.05%, 胜率 52.5% +- 美股涨0.5~1.5%(1180天): A股跳空 N=1180, 均值 0.08%, 中位数 0.07%, 胜率 61.4%; 日内 N=1180, 均值 0.14%, 中位数 0.12%, 胜率 55.7%; 全天 N=1180, 均值 0.23%, 中位数 0.17%, 胜率 59.8% +- 美股涨>1.5%(287天): A股跳空 N=287, 均值 0.57%, 中位数 0.37%, 胜率 84.0%; 日内 N=287, 均值 -0.11%, 中位数 -0.14%, 胜率 44.9%; 全天 N=287, 均值 0.45%, 中位数 0.28%, 胜率 60.3% + +### 道琼斯T日 -> 沪深300T+1日 + +- 美股跌>1.5%(324天): A股跳空 N=324, 均值 -1.01%, 中位数 -0.83%, 胜率 9.6%; 日内 N=324, 均值 0.39%, 中位数 0.20%, 胜率 57.7%; 全天 N=324, 均值 -0.63%, 中位数 -0.49%, 胜率 37.7% +- 美股跌0.5~1.5%(890天): A股跳空 N=890, 均值 -0.26%, 中位数 -0.23%, 胜率 24.3%; 日内 N=890, 均值 0.13%, 中位数 0.04%, 胜率 51.9%; 全天 N=890, 均值 -0.13%, 中位数 -0.10%, 胜率 45.2% +- 美股正负0.5%内(2894天): A股跳空 N=2894, 均值 -0.04%, 中位数 -0.04%, 胜率 44.7%; 日内 N=2894, 均值 0.07%, 中位数 0.01%, 胜率 50.4%; 全天 N=2894, 均值 0.03%, 中位数 0.03%, 胜率 51.5% +- 美股涨0.5~1.5%(1180天): A股跳空 N=1180, 均值 0.12%, 中位数 0.11%, 胜率 64.8%; 日内 N=1180, 均值 0.13%, 中位数 0.04%, 胜率 51.2%; 全天 N=1180, 均值 0.25%, 中位数 0.17%, 胜率 58.1% +- 美股涨>1.5%(287天): A股跳空 N=287, 均值 0.64%, 中位数 0.45%, 胜率 83.6%; 日内 N=287, 均值 -0.19%, 中位数 -0.19%, 胜率 42.2%; 全天 N=287, 均值 0.45%, 中位数 0.32%, 胜率 61.0% diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..5340e2e --- /dev/null +++ b/requirements.txt @@ -0,0 +1,8 @@ +akshare>=1.18.92 +pandas>=2.0 +aiohttp>=3.9 +lxml +beautifulsoup4 +py_mini_racer +tqdm +requests diff --git a/src/realtime/collector.py b/src/realtime/collector.py new file mode 100644 index 0000000..23febff --- /dev/null +++ b/src/realtime/collector.py @@ -0,0 +1,267 @@ +# -*- coding: utf-8 -*- +""" +时间线采集引擎(跨平台):将 realtime_analyzer 的分析循环重构为可嵌入的采集器。 +- 盘前:美股隔夜收盘 -> 实证先验(a-stock-timeline-patterns skill)输出今日预估 +- 盘中:东财 7x24 快讯增量入库 + 关键词告警 + 上证实时点位 +- 盘后:龙虎榜入库 + 收盘日报 +事件通过回调 emit(event: dict) 交给上层(桌面版写文件,B/S 版走 WebSocket 广播)。 +""" +import json +import os +import sqlite3 +import traceback +from datetime import datetime, timedelta +from pathlib import Path + +import requests + +ROOT = Path(__file__).resolve().parent.parent.parent +DB = ROOT / 'data' / 'a_stock.db' +FEED = ROOT / 'data' / 'realtime_feed.jsonl' +SKILL_REF = Path(r'E:\Data\skills\a-stock-timeline-patterns\references\overnight_transmission.md') + +KEYWORDS = ['降息', '降准', '加息', '关税', '制裁', '收购', '重组', '国债', '证监会', 'PMI', 'CPI'] + +# 出站请求域名白名单(SSRF 防护) +_ALLOWED_HOSTS = {'np-weblist.eastmoney.com', 'np-listapi.eastmoney.com'} + + +def _safe_get(url, params=None, timeout=10): + from urllib.parse import urlparse + u = urlparse(url) + if u.scheme != 'https' or u.hostname not in _ALLOWED_HOSTS: + raise ValueError('blocked non-allowlist url: %s' % url) + return requests.get(url, params=params, timeout=timeout, + headers={'User-Agent': 'Mozilla/5.0'}, allow_redirects=False) + + +class TimelineCollector: + """常驻采集引擎:一个线程安全的同步循环,由上层调度(线程/异步 executor)""" + + def __init__(self, emit, db_path=None): + self.emit = emit # callable(dict) + self.db_path = str(db_path or DB) + self.state = {'date': datetime.now().strftime('%Y-%m-%d')} + self.priors = {} + self._load_priors() + + # ---------- 基础 ---------- + + def _conn(self): + conn = sqlite3.connect(self.db_path, check_same_thread=False) + return conn + + def _emit(self, kind, data): + rec = {'ts': datetime.now().strftime('%Y-%m-%d %H:%M:%S'), 'kind': kind, 'data': data} + try: + FEED.parent.mkdir(parents=True, exist_ok=True) + with open(FEED, 'a', encoding='utf-8') as f: + f.write(json.dumps(rec, ensure_ascii=False) + '\n') + except OSError: + pass + self.emit(rec) + + # ---------- 实证先验(skill 注入) ---------- + + def _load_priors(self): + import re + self.priors = {} + path = str(SKILL_REF) + if not os.path.exists(path): + # Linux 部署:skill 文件可放项目 docs/ 下 + alt = Path(__file__).resolve().parent.parent.parent / 'docs' / 'overnight_transmission.md' + if not alt.exists(): + return + path = str(alt) + try: + with open(path, encoding='utf-8') as f: + cur_us = cur_idx = None + for line in f: + m = re.match(r'### (\S+)T日 -> (\S+)T\+1日', line) + if m: + cur_us, cur_idx = m.group(1), m.group(2) + continue + m = re.match( + r'- 美股(\S+?)((\d+)天): A股跳空 N=\d+, 均值 (-?[\d.]+)%, 中位数 (-?[\d.]+)%, ' + r'胜率 ([\d.]+)%; 日内 N=\d+, 均值 (-?[\d.]+)%, 中位数 (-?[\d.]+)%, 胜率 ([\d.]+)%', + line) + if m and cur_us and cur_idx: + self.priors[(cur_us, cur_idx, m.group(1))] = { + 'n': int(m.group(2)), + 'gap_mean': float(m.group(3)), 'gap_median': float(m.group(4)), + 'gap_win': float(m.group(5)), + 'intra_mean': float(m.group(6)), 'intra_median': float(m.group(7)), + 'intra_win': float(m.group(8)), + } + except OSError: + pass + + # ---------- 数据获取 ---------- + + def us_overnight(self, conn): + out = {} + for code in ('US.NDX', 'US.SPX', 'US.DJI'): + row = conn.execute( + "SELECT trade_date, pct_change FROM global_index WHERE index_code=? " + "ORDER BY trade_date DESC LIMIT 1", (code,)).fetchone() + if row: + out[code] = {'date': row[0], 'pct': row[1] or 0.0} + return out + + def forecast_from_priors(self, us): + if not self.priors or 'US.NDX' not in us: + return None + pct = us['US.NDX']['pct'] / 100.0 + if pct <= -0.015: + b = '跌>1.5%' + elif pct <= -0.005: + b = '跌0.5~1.5%' + elif pct < 0.005: + b = '正负0.5%内' + elif pct < 0.015: + b = '涨0.5~1.5%' + else: + b = '涨>1.5%' + k = ('US.NDX', 'sh000001', b) + if k not in self.priors: + return None + p = self.priors[k] + return ('隔夜预估[美股纳指{:+.2f}% -> 分箱"{}"]: 历史上上证次日跳空均值 {:+.2f}%' + '(低开概率 {:.0f}%),日内均值 {:+.2f}%(日内收涨概率 {:.0f}%),全天均值 {:+.2f}%。' + '(样本{}天,美股先验,仅参考)').format( + us['US.NDX']['pct'], b, p['gap_mean'], 100 - p['gap_win'], + p['intra_mean'], p['intra_win'], p['gap_mean'] + p['intra_mean'], p['n']) + + def index_spot_sh(self): + try: + import akshare as ak + df = ak.stock_zh_index_spot_em(symbol='上证系列指数') + row = df[df['名称'] == '上证指数'] + if not row.empty: + r = row.iloc[0] + return {'price': float(r['最新价']), 'pct': float(r['涨跌幅'])} + except Exception: + pass + try: + import akshare as ak + df = ak.stock_zh_index_spot_sina() + row = df[df['代码'] == 'sh000001'] + if not row.empty: + r = row.iloc[0] + return {'price': float(r['最新价']), 'pct': float(r['涨跌幅'])} + except Exception: + pass + return None + + def fetch_em_flash(self, conn): + """东财 7x24 快讯增量入库,返回 (time, text) 新增列表""" + fresh = [] + try: + r = _safe_get('https://np-weblist.eastmoney.com/comm/web/getFastNewsList', + params={'client': 'web', 'biz': 'web_724', 'fastColumn': '102', + 'sortEnd': '', 'pageSize': '20', 'req_trace': '1'}) + data = r.json().get('data', {}) or {} + for n in data.get('fastNewsList', []) or []: + ts = n.get('showTime', '') + summary = (n.get('summary') or n.get('title') or '').strip() + if not ts or not summary: + continue + dup = conn.execute( + "SELECT 1 FROM news_flash WHERE event_time=? AND content=? LIMIT 1", + (ts, summary)).fetchone() + if not dup: + conn.execute( + "INSERT INTO news_flash(content, event_time, source, inserted_at) " + "VALUES (?,?,?,datetime('now','localtime'))", (summary, ts, '东财7x24')) + fresh.append((ts, summary[:60])) + conn.commit() + except Exception as e: + self._emit('采集错误', {'msg': str(e)[:120]}) + return fresh + + def save_lhb_today(self, conn): + try: + from src.fetcher.lhb_fetcher import fetch_lhb + from src.storage.db import upsert_rows + df = fetch_lhb() + if df is not None and not df.empty: + return upsert_rows(df, 'lhb_daily', + conflict_cols=['trade_date', 'ts_code', 'reason']) + except Exception as e: + self._emit('采集错误', {'msg': str(e)[:120]}) + return 0 + + def news_alerts(self, conn): + hits = [] + cutoff = (datetime.now() - timedelta(minutes=30)).strftime('%Y-%m-%d %H:%M:%S') + for table, tcol, ccol in (('news_flash', 'event_time', 'content'), + ('news_cn', 'pub_date', 'title')): + try: + rows = conn.execute( + "SELECT {}, {} FROM {} WHERE {} >= ? ORDER BY {} DESC LIMIT 50".format( + tcol, ccol, table, tcol, tcol), (cutoff,)).fetchall() + except sqlite3.OperationalError: + continue + for ts, text in rows: + for kw in KEYWORDS: + if kw in (text or ''): + hits.append({'time': ts, 'kw': kw, 'text': (text or '')[:80]}) + break + return hits + + # ---------- 单轮采集 ---------- + + def one_cycle(self): + conn = self._conn() + try: + now = datetime.now() + hm = now.hour * 100 + now.minute + + if 700 <= hm < 925 and not self.state.get('premarket_done'): + us = self.us_overnight(conn) + fc = self.forecast_from_priors(us) + self._emit('盘前隔夜预估', {'us': us, 'forecast': fc}) + self.state['premarket_done'] = True + + if 925 <= hm < 1505: + for ts, txt in self.fetch_em_flash(conn): + self._emit('新快讯', {'time': ts, 'text': txt}) + spot = self.index_spot_sh() + if spot and spot.get('price', 0) > 0: + self._emit('盘中点位', {'上证': spot['price'], '涨跌幅%': spot['pct']}) + for h in self.news_alerts(conn): + self._emit('快讯关键词告警', h) + + if hm >= 1510 and not self.state.get('postmarket_done'): + n = self.save_lhb_today(conn) + sh = conn.execute("SELECT trade_date, close FROM stock_daily " + "WHERE ts_code='sh000001' ORDER BY trade_date DESC LIMIT 1").fetchone() + self._emit('收盘日报', {'龙虎榜新增': n, '上证最新收盘': sh}) + self.state['postmarket_done'] = True + + if self.state.get('date') != now.strftime('%Y-%m-%d'): + self.state.clear() + self.state['date'] = now.strftime('%Y-%m-%d') + self._load_priors() + self._emit('日切', {'date': self.state['date']}) + finally: + conn.close() + + def run_forever(self, interval=60, on_error=None): + """阻塞式常驻循环(Linux 服务/桌面版均可直接调用)""" + while True: + try: + self.one_cycle() + except Exception: + traceback.print_exc() + if on_error: + on_error() + time.sleep(interval) + + +if __name__ == '__main__': + c = TimelineCollector(emit=lambda rec: print( + '[{}] {} {}'.format(rec['ts'], rec['kind'], + json.dumps(rec['data'], ensure_ascii=False)), flush=True)) + print('时间线采集引擎启动(Ctrl+C 停止)', flush=True) + c.run_forever() diff --git a/src/web/index.html b/src/web/index.html new file mode 100644 index 0000000..85e9a84 --- /dev/null +++ b/src/web/index.html @@ -0,0 +1,125 @@ + + + + + +JQuant · A股时间线实时监控 + + + +
+

JQ · A股时间线实时监控 B/S(Linux 服务端推送)

+ 未连接 +
+
服务端:抓取(TDX/东财/同花顺)+ 分析(波动分解/事件归因)每 60s 一轮 → WebSocket 实时推送本页 | 缓冲事件 0
+
+
+ + + + + + +
+ +
+ + + diff --git a/src/web/server.py b/src/web/server.py new file mode 100644 index 0000000..59b558f --- /dev/null +++ b/src/web/server.py @@ -0,0 +1,110 @@ +# -*- coding: utf-8 -*- +""" +B/S 服务端:时间线采集引擎(后台线程)+ WebSocket 实时推送 + Web 仪表盘 +Linux 服务器常驻运行: python -m src.web.server (或 systemd,见 deploy/a-stock-timeline.service) +端口默认 8100,可用环境变量 PORT 覆盖。 +""" +import asyncio +import json +import os +from collections import deque +from pathlib import Path + +from aiohttp import WSMsgType, web + +import sys +sys.path.insert(0, str(Path(__file__).resolve().parent.parent.parent)) + +from src.realtime.collector import TimelineCollector # noqa: E402 + +PORT = int(os.environ.get('PORT', '8100')) +ROOT = Path(__file__).resolve().parent.parent.parent +WEB_DIR = Path(__file__).resolve().parent +RECENT_MAX = 500 + +recent = deque(maxlen=RECENT_MAX) # 最近事件(内存) +clients = set() # 活跃 WS 连接 + + +async def broadcast(event: dict): + recent.append(event) + payload = json.dumps(event, ensure_ascii=False) + dead = set() + for ws in clients: + try: + await ws.send_str(payload) + except Exception: + dead.add(ws) + for ws in dead: + clients.discard(ws) + + +async def collector_task(_app): + """把阻塞式采集循环放进线程池,事件桥接到 asyncio 广播""" + loop = asyncio.get_running_loop() + collector = TimelineCollector(emit=lambda rec: loop.call_soon_threadsafe( + asyncio.ensure_future, broadcast(rec))) + + async def poll(): + while True: + await loop.run_in_executor(None, collector.one_cycle) + await asyncio.sleep(60) + + task = asyncio.create_task(poll()) + # 启动即推最近历史(从 jsonl 恢复) + feed = ROOT / 'data' / 'realtime_feed.jsonl' + if feed.exists(): + try: + lines = feed.read_text(encoding='utf-8').strip().splitlines()[-200:] + for line in lines: + try: + recent.append(json.loads(line)) + except json.JSONDecodeError: + pass + except OSError: + pass + yield + task.cancel() + + +async def index(_request): + return web.FileResponse(WEB_DIR / 'index.html') + + +async def recent_events(_request): + return web.json_response({'events': list(recent)}) + + +async def stats(_request): + return web.json_response({'clients': len(clients), 'buffered': len(recent)}) + + +async def ws_handler(request): + ws = web.WebSocketResponse(heartbeat=30) + await ws.prepare(request) + clients.add(ws) + await ws.send_str(json.dumps({'kind': '连接成功', + 'data': {'msg': '实时推送已连接', 'buffered': len(recent)}}, + ensure_ascii=False)) + try: + async for msg in ws: + if msg.type == WSMsgType.ERROR: + break + finally: + clients.discard(ws) + return ws + + +def build_app(): + app = web.Application() + app.router.add_get('/', index) + app.router.add_get('/api/recent', recent_events) + app.router.add_get('/api/stats', stats) + app.router.add_get('/ws', ws_handler) + app.cleanup_ctx.append(collector_task) + return app + + +if __name__ == '__main__': + print('JQuant 时间线 B/S 服务启动: http://0.0.0.0:{} (WS: /ws)'.format(PORT)) + web.run_app(build_app(), host='0.0.0.0', port=PORT) diff --git a/start_server.sh b/start_server.sh new file mode 100644 index 0000000..f7d7b24 --- /dev/null +++ b/start_server.sh @@ -0,0 +1,3 @@ +#!/usr/bin/env bash +cd "$(dirname "$0")" +exec python3 -m src.web.server