diff --git a/research/live/compare_sites.py b/research/live/compare_sites.py new file mode 100644 index 0000000..8bf7add --- /dev/null +++ b/research/live/compare_sites.py @@ -0,0 +1,184 @@ +"""对比两个(或多个)采集站点的数据差异。 + +部署第二台机器的目的就是这个:延迟是「本地接收 − 交易所 K 线收盘」,直接 +取决于机器到交易所的网络距离,换个地理位置这个数会变。而滑点里最大的一项 +是延迟漂移,所以站点选址本身就是一个可优化的参数。 + +## 判读前必须先看的两件事 + +1. **时钟。** 两台机器的时钟偏移差多少,延迟对比就凭空差多少,且不报错。 + run_meta_*.json 里有各站启动时的 chrony 偏移,先确认都在 10ms 内。 +2. **同期。** 只比两站都有数据的那些 kline_ts。不取交集的话,比的可能是 + 不同时段的市场状态,而延迟对市场活跃度是敏感的。 + +## 用法 + +把各站的 research/out/ 收到一处(文件名相同会覆盖,所以先按站点改名或 +分目录放),然后: + + python research/live/compare_sites.py --glob 'collected/*/shadow_latency.csv' +""" +from __future__ import annotations + +import argparse +import glob +import json +from pathlib import Path + +import numpy as np +import pandas as pd + + +def load(patterns: list[str]) -> pd.DataFrame: + paths: list[str] = [] + for p in patterns: + paths.extend(sorted(glob.glob(p))) + if not paths: + raise SystemExit(f"没有匹配到文件:{patterns}") + frames = [] + for p in paths: + df = pd.read_csv(p) + if "site" not in df.columns: + raise SystemExit( + f"{p} 没有 site 列。这是 2026-08-28 之前采的旧数据," + f"无法确定来源,不能用于跨地对比") + df["_src"] = p + frames.append(df) + out = pd.concat(frames, ignore_index=True) + print(f"读入 {len(paths)} 个文件、{len(out):,} 行、" + f"站点 {sorted(out['site'].unique())}") + return out + + +def show_meta(out_dirs: list[Path]) -> None: + print("\n########## 一、运行元数据 ##########") + metas = [] + for d in out_dirs: + metas.extend(sorted(d.glob("run_meta_*.json"))) + if not metas: + print(" 没找到 run_meta_*.json。时钟偏移与代码版本无法核对——") + print(" 两站数据若有差异,分不清是地理位置还是环境不同造成的") + return + rows = [] + for m in metas: + try: + rows.append(json.loads(m.read_text())) + except Exception as e: + print(f" {m.name} 读取失败:{e!r}") + if not rows: + return + df = pd.DataFrame(rows) + keep = [c for c in ("site", "clock_offset_ms", "git_commit", "git_dirty", + "image_digest", "nproc", "mem_gb", "tz", + "started_utc") if c in df.columns] + print(df[keep].to_string(index=False)) + if "clock_offset_ms" in df and df["clock_offset_ms"].notna().any(): + o = df["clock_offset_ms"].astype(float) + spread = float(o.max() - o.min()) + flag = "" if spread < 5 else " ⚠ 这个差会直接叠加到延迟对比上" + print(f"\n 站点间时钟偏移极差 {spread:.3f}ms{flag}") + if "git_commit" in df and df["git_commit"].nunique() > 1: + print(" ⚠ 各站代码版本不同,差异可能来自代码而非地理位置") + if "git_dirty" in df and df["git_dirty"].any(): + print(" ⚠ 有站点带未提交改动,无法复现") + + +def compare_latency(df: pd.DataFrame, col: str = "lag_data_ms") -> None: + """延迟对比。只取各站都有的 kline_ts,避免比到不同时段。""" + if col not in df.columns: + print(f"\n没有 {col} 列") + return + sites = sorted(df["site"].unique()) + if len(sites) < 2: + print(f"\n只有一个站点({sites[0]}),无从对比。" + f"等第二台机器的数据到齐") + return + + print(f"\n########## 二、到达延迟({col}) ##########") + print("\n 全量(各站各自的样本,时段可能不同)") + for s in sites: + x = df[df["site"] == s][col].dropna().astype(float) + print(f" {s:<16} n={len(x):>6} 中位 {x.median():>7.0f}ms " + f"P90 {np.percentile(x, 90):>7.0f}ms " + f"P99 {np.percentile(x, 99):>7.0f}ms") + + # 取交集:同一根 K 线在各站都有记录 + key = ["sym", "kline_ts"] + piv = df.pivot_table(index=key, columns="site", values=col, + aggfunc="first") + both = piv.dropna() + if both.empty: + print("\n 各站没有共同的 K 线。可能是采集时段不重叠,") + print(" 或 kline_ts 对不上(先查两站时区与时钟)") + return + print(f"\n 同根对比({len(both):,} 根 K 线,各站都有)") + for s in sites: + x = both[s].astype(float) + print(f" {s:<16} 中位 {x.median():>7.0f}ms " + f"P90 {np.percentile(x, 90):>7.0f}ms") + base = sites[0] + for s in sites[1:]: + d = (both[s] - both[base]).astype(float) + # 配对差的符号检验:同根配对消掉了市场状态,比两个中位数相减干净 + n_pos = int((d > 0).sum()) + print(f"\n {s} − {base}:中位差 {d.median():+.0f}ms " + f"· 均值差 {d.mean():+.0f}ms") + print(f" {s} 更慢的根占 {n_pos / len(d) * 100:.1f}%" + f"(50% 表示无系统性差异)") + for sym in sorted(both.index.get_level_values("sym").unique()): + ds = d.xs(sym, level="sym") + print(f" {sym:<5} 中位差 {ds.median():+7.0f}ms (n={len(ds)})") + + +def compare_drift(patterns: list[str]) -> None: + """漂移对比。延迟差若能兑换成漂移差,才是钱上的差别。""" + paths: list[str] = [] + for p in patterns: + paths.extend(sorted(glob.glob(p))) + if not paths: + return + frames = [] + for p in paths: + d = pd.read_csv(p) + if "site" in d.columns: + frames.append(d) + if not frames: + return + df = pd.concat(frames, ignore_index=True) + if df["site"].nunique() < 2: + return + print("\n########## 三、延迟漂移(无条件,每根都记) ##########") + for label in sorted(df["delay_label"].dropna().unique()): + sub = df[df["delay_label"] == label] + line = f" {label:>7}" + for s in sorted(sub["site"].unique()): + x = sub[sub["site"] == s]["drift_bp_long"].dropna().astype(float) + if len(x): + line += f" · {s} {x.abs().median():.3f}bp(n={len(x)})" + print(line) + print("\n 同一个固定延迟点上,两站的漂移应当几乎相同——漂移是市场性质,") + print(" 与机器位置无关。若差异明显,先查时钟与采集时段是否对齐") + + +def main() -> None: + ap = argparse.ArgumentParser() + ap.add_argument("--glob", action="append", default=None, + help="shadow_latency.csv 的路径模式,可给多次") + ap.add_argument("--drift-glob", action="append", default=None) + ap.add_argument("--meta-dir", action="append", default=None) + a = ap.parse_args() + + lat = a.glob or ["research/out/shadow_latency.csv", + "collected/*/shadow_latency.csv"] + drf = a.drift_glob or ["research/out/shadow_drift.csv", + "collected/*/shadow_drift.csv"] + metas = [Path(p) for p in (a.meta_dir or ["research/out", "collected"])] + + show_meta([p for p in metas if p.is_dir()]) + df = load(lat) + compare_latency(df) + compare_drift(drf) + + +if __name__ == "__main__": + main() diff --git a/research/live/deploy/README.md b/research/live/deploy/README.md new file mode 100644 index 0000000..898ccee --- /dev/null +++ b/research/live/deploy/README.md @@ -0,0 +1,97 @@ +# 影子采集器部署 + +在第二台机器上跑一套完全相同的采集,用来看不同地理位置的数据差异。 + +## 为什么要跨地采集 + +滑点里最大的一项是**延迟漂移**:从 K 线收盘到实际下单之间,价格已经走掉的 +那部分。而延迟 = 交易所出包 + 网络传输 + 本地处理。本机(新加坡)实测到达 +延迟中位 350~650ms,其中网络传输占多少、换个机房能压掉多少,只有实测。 + +延迟压下来直接等于滑点下降,所以机房选址是个可优化参数,不是固定成本。 + +## 三个必须一致,一个必须不同 + +必须一致,否则差异分不清是地理位置还是环境造成的: + +- **代码版本**(`git_commit`)——同一个 commit +- **镜像摘要**(`image_digest`)——同一个 hummingbot 镜像 +- **时钟**——两台都同步到 NTP,偏移都在 10ms 内 + +必须不同: + +- **`SHADOW_SITE`**——写进每一行数据,是合并后区分来源的唯一依据 + +`start.sh` 会把这四项连同内核、核数、内存一起写进 +`research/out/run_meta_.json`。两地数据对不上时先看这个文件。 + +## 时钟为什么是硬门槛 + +所有延迟数字都是「本地时钟 − 交易所 K 线收盘时间戳」。时钟偏 50ms,全部 +延迟就同向偏 50ms,而且**不会有任何报错**——只会让跨地对比得出一个干净、 +自信、且完全错误的结论。所以 `start.sh` 在时钟未同步或偏移超阈值时直接 +拒绝启动,而不是打个警告了事。 + +## 步骤 + +在新机器上: + +```bash +git clone ssh://jack@git.jackyu66.com:2222/jack/chan.git +cd chan && git checkout chan + +bash research/live/deploy/setup.sh # 装 docker + chrony,拉镜像 +SHADOW_SITE=aws-tokyo bash research/live/deploy/start.sh +``` + +确认健康: + +```bash +bash research/live/deploy/status.sh +``` + +启动日志里应当看到: + +``` +[补丁] 覆盖生效:基类取首元素 … 本地取末元素 … +成交流已挂 ['BTC', 'ETH', 'SOL'] +就绪 3.0s · 1m [2001, 2001, 2001] 根 · 5m [801, 801, 801] 根 +``` + +第一行尤其重要。上游 Bitget 连接器的换根解析有 bug(只取多根消息的首元素), +补丁把它修掉拿回约 1.06 秒。补丁若失效是静默的——不崩不报错,只是延迟悄悄 +退回 1.4 秒,所以启动时做了断言。 + +## 对比 + +把两站的 `research/out/` 收到一处(同名文件会覆盖,所以分目录放): + +```bash +mkdir -p collected/sg collected/aws +rsync -av sg-box:chan/research/out/ collected/sg/ +rsync -av aws-box:chan/research/out/ collected/aws/ + +python research/live/compare_sites.py \ + --glob 'collected/*/shadow_latency.csv' \ + --drift-glob 'collected/*/shadow_drift.csv' \ + --meta-dir collected/sg --meta-dir collected/aws +``` + +对比脚本做两件事值得说明: + +- **只取各站都有的 K 线**做配对比较。不取交集就可能在比不同时段,而延迟对 + 市场活跃度敏感。 +- 报**配对差的符号占比**而不只是两个中位数相减。同根配对消掉了市场状态, + 「A 比 B 慢的根占多少」比「两个中位数差多少」更能说明有无系统性差异。 + +判读上有一条自检:**同一个固定延迟点上,两站的漂移应当几乎相同**——漂移是 +市场性质,与机器位置无关。若漂移也差很多,先怀疑时钟或时段没对齐,而不是 +急着下结论。 + +## 资源占用 + +本机实测:内存约 1.5GB(两个计算进程 + 盘口缓冲),CPU 单核不满。 +盘口与成交流落盘约 15MB/天(gzip)。一周 168 小时的量级在百 MB 内。 + +`--workers 2` 是因为信号计算走独立进程池、不能阻塞事件循环。核数少的机型 +可以给 1,但要看心跳里的 `compute_ms`:若接近 60 秒就会开始堆积。 diff --git a/research/live/deploy/setup.sh b/research/live/deploy/setup.sh new file mode 100755 index 0000000..305ae2f --- /dev/null +++ b/research/live/deploy/setup.sh @@ -0,0 +1,87 @@ +#!/usr/bin/env bash +# 影子采集器的机器初始化。幂等,可重复跑。 +# +# 用途:在另一台机器(如 AWS)上部署一套完全相同的采集,用来看不同地理位置 +# 的数据采集有无差异。要比的主要是延迟——「本地接收 − K线收盘」这个量直接 +# 取决于机器到交易所的网络距离。 +# +# 时钟同步是硬前置条件,不是可选项。所有延迟数字都是本地时钟减交易所时间戳, +# 时钟偏 50ms 就等于所有延迟凭空多(或少)50ms,而且不会有任何报错。所以这里 +# 装并启用 NTP,start.sh 里还会再校验一次、不合格拒绝启动。 +# +# bash research/live/deploy/setup.sh +set -euo pipefail + +IMAGE="${SHADOW_IMAGE:-hummingbot/hummingbot:latest}" + +say() { printf '\n\033[1m==> %s\033[0m\n' "$*"; } + +say "系统信息" +uname -a +echo "内存:$(free -g | awk '/^Mem:/{print $2"GB 总 / "$7"GB 可用"}')" +echo "CPU:$(nproc) 核" + +say "安装 docker" +if command -v docker >/dev/null 2>&1; then + echo "已有 docker $(docker --version)" +else + if command -v apt-get >/dev/null 2>&1; then + sudo apt-get update -qq + sudo apt-get install -y -qq ca-certificates curl gnupg + sudo install -m 0755 -d /etc/apt/keyrings + curl -fsSL https://download.docker.com/linux/ubuntu/gpg \ + | sudo gpg --dearmor -o /etc/apt/keyrings/docker.gpg + sudo chmod a+r /etc/apt/keyrings/docker.gpg + . /etc/os-release + echo "deb [arch=$(dpkg --print-architecture) signed-by=/etc/apt/keyrings/docker.gpg] \ +https://download.docker.com/linux/ubuntu ${VERSION_CODENAME} stable" \ + | sudo tee /etc/apt/sources.list.d/docker.list >/dev/null + sudo apt-get update -qq + sudo apt-get install -y -qq docker-ce docker-ce-cli containerd.io + elif command -v dnf >/dev/null 2>&1; then + # Amazon Linux 2023 + sudo dnf install -y -q docker + sudo systemctl enable --now docker + else + echo "不认识的包管理器,请手动装 docker" >&2 + exit 1 + fi + sudo usermod -aG docker "$USER" || true + echo "已装 docker。若本次 shell 无权限,重新登录后再跑 start.sh" +fi + +say "启用时钟同步(延迟测量的前置条件)" +if command -v chronyc >/dev/null 2>&1; then + echo "已有 chrony" +elif command -v apt-get >/dev/null 2>&1; then + sudo apt-get install -y -qq chrony +elif command -v dnf >/dev/null 2>&1; then + sudo dnf install -y -q chrony +fi +sudo systemctl enable --now chrony 2>/dev/null \ + || sudo systemctl enable --now chronyd 2>/dev/null || true +sleep 3 +if command -v chronyc >/dev/null 2>&1; then + chronyc tracking | grep -E 'Reference ID|System time|Last offset' || true +else + timedatectl 2>/dev/null | grep -i synchron || true +fi + +say "拉镜像 $IMAGE" +docker pull "$IMAGE" +docker image inspect "$IMAGE" --format '摘要 {{index .RepoDigests 0}}' 2>/dev/null || true + +say "准备输出目录" +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../../.." && pwd)" +mkdir -p "$REPO_ROOT/research/out" +echo "$REPO_ROOT/research/out" + +say "完成" +cat <<'EOF' +下一步: + + SHADOW_SITE=aws-tokyo bash research/live/deploy/start.sh + +SHADOW_SITE 必须显式给且两台机器不能相同——它会写进每一行数据, +是之后区分数据来源的唯一依据。 +EOF diff --git a/research/live/deploy/start.sh b/research/live/deploy/start.sh new file mode 100755 index 0000000..42f7803 --- /dev/null +++ b/research/live/deploy/start.sh @@ -0,0 +1,115 @@ +#!/usr/bin/env bash +# 启动影子采集器。跑之前先 setup.sh。 +# +# SHADOW_SITE=aws-tokyo bash research/live/deploy/start.sh +# SHADOW_SITE=aws-tokyo HOURS=168 WORKERS=2 bash research/live/deploy/start.sh +# +# 为什么 SHADOW_SITE 必填:它写进每一行数据,是两台机器的数据合起来之后 +# 唯一的来源区分。缺了就只能靠文件路径猜,一合并就分不清了。 +# +# 为什么时钟不同步就拒绝启动:所有延迟数字都是「本地时钟 − 交易所 K 线收盘 +# 时间戳」。时钟偏 50ms,全部延迟就凭空偏 50ms,而这**不会有任何报错**, +# 只会让跨地对比得出一个完全错误的结论。这类静默错误必须在启动就挡掉。 +set -euo pipefail + +NAME="${NAME:-shadow}" +IMAGE="${SHADOW_IMAGE:-hummingbot/hummingbot:latest}" +HOURS="${HOURS:-168}" +WORKERS="${WORKERS:-2}" +MAX_OFFSET_MS="${MAX_OFFSET_MS:-10}" + +REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../../.." && pwd)" +OUT="$REPO_ROOT/research/out" + +die() { printf '\033[31m错误:%s\033[0m\n' "$*" >&2; exit 1; } +say() { printf '\n\033[1m==> %s\033[0m\n' "$*"; } + +[[ -n "${SHADOW_SITE:-}" ]] || die "必须设 SHADOW_SITE,例如 SHADOW_SITE=aws-tokyo。 +它写进每一行数据,是跨地对比时区分来源的唯一依据。" + +say "站点 $SHADOW_SITE" + +say "校验时钟同步" +offset_ms="" +if command -v chronyc >/dev/null 2>&1; then + if ! chronyc tracking >/dev/null 2>&1; then + die "chrony 没在跑。先 sudo systemctl start chrony(或 chronyd)" + fi + # System time 那行形如 "0.000058703 seconds fast of NTP time" + line="$(chronyc tracking | grep '^System time' || true)" + secs="$(awk '{print $4}' <<<"$line")" + offset_ms="$(awk -v s="$secs" 'BEGIN{printf "%.3f", s*1000}')" + echo "$line" + leap="$(chronyc tracking | awk -F': *' '/Leap status/{print $2}')" + [[ "$leap" == "Normal" ]] || die "chrony leap status = $leap,尚未收敛。等几分钟再试" + over="$(awk -v o="$offset_ms" -v m="$MAX_OFFSET_MS" 'BEGIN{print (o>m)?1:0}')" + [[ "$over" == "0" ]] || die "时钟偏移 ${offset_ms}ms 超过阈值 ${MAX_OFFSET_MS}ms。 +所有延迟测量都会同向偏这么多且不报错,跨地对比会得出错误结论。 +先等 chrony 收敛,或调 MAX_OFFSET_MS(不建议)。" + echo "偏移 ${offset_ms}ms,在 ${MAX_OFFSET_MS}ms 阈值内" +elif command -v timedatectl >/dev/null 2>&1; then + timedatectl | grep -qi 'synchronized: yes' \ + || die "系统时钟未同步。装 chrony:见 setup.sh" + echo "timedatectl 报已同步(无 chronyc,拿不到具体偏移)" +else + die "既无 chronyc 也无 timedatectl,无法确认时钟。装 chrony 后再启动" +fi + +say "检查镜像" +docker image inspect "$IMAGE" >/dev/null 2>&1 || die "没有镜像 $IMAGE,先跑 setup.sh" +digest="$(docker image inspect "$IMAGE" --format '{{if .RepoDigests}}{{index .RepoDigests 0}}{{end}}' 2>/dev/null || true)" + +say "停掉旧容器" +docker rm -f "$NAME" >/dev/null 2>&1 || true + +mkdir -p "$OUT" + +# 运行元数据。两地数据对不上时,先看这个文件——镜像摘要、代码版本、时钟偏移 +# 三者任一不同都足以解释差异,不必去猜。 +meta="$OUT/run_meta_${SHADOW_SITE}.json" +cat >"$meta" </dev/null || echo unknown)", + "git_dirty": $(git -C "$REPO_ROOT" diff --quiet 2>/dev/null && echo false || echo true), + "image": "$IMAGE", + "image_digest": "${digest:-unknown}", + "hours": $HOURS, + "workers": $WORKERS, + "kernel": "$(uname -r)", + "nproc": $(nproc), + "mem_gb": $(free -g | awk '/^Mem:/{print $2}') +} +EOF +echo "元数据已写 $meta" + +say "启动容器 $NAME" +docker run -d --name "$NAME" -w /home/hummingbot \ + --restart unless-stopped \ + -e PYTHONPATH=/home/hummingbot:/repo/research:/repo/research/live:/repo \ + -e SHADOW_SITE="$SHADOW_SITE" \ + -v "$REPO_ROOT:/repo:ro" \ + -v "$OUT:/out" \ + --entrypoint /opt/conda/envs/hummingbot/bin/python \ + "$IMAGE" /repo/research/live/shadow_hb.py \ + --hours "$HOURS" --workers "$WORKERS" >/dev/null + +echo "已启动。等启动自检(补丁断言 + 历史回填,约 60 秒)…" +sleep 45 +docker logs "$NAME" 2>&1 | tail -12 + +cat < %s\033[0m\n' "$*"; } + +say "容器" +docker ps -a --filter "name=^${NAME}$" \ + --format 'table {{.Names}}\t{{.Status}}\t{{.RunningFor}}' || true + +say "时钟" +if command -v chronyc >/dev/null 2>&1; then + chronyc tracking | grep -E 'System time|Last offset|Leap status' +fi + +say "最近心跳" +docker logs "$NAME" 2>&1 | grep '\[心跳\]' | tail -3 || echo "还没到第一次心跳(每 5 分钟一次)" + +say "告警与异常" +docker logs "$NAME" 2>&1 \ + | grep -E '⚠|错误|失效|损坏|停滞|Traceback|退化' | tail -10 \ + || echo "无" + +say "输出规模" +for f in shadow_latency.csv shadow_drift.csv shadow_signals.csv; do + p="$OUT/$f" + if [[ -s "$p" ]]; then + printf ' %-22s %8d 行\n' "$f" "$(( $(wc -l <"$p") - 1 ))" + elif [[ -f "$p" ]]; then + # 已建但表头还没冲刷:signals 只在有信号时才 flush,门控后每天仅 4~6 个 + printf ' %-22s %8s\n' "$f" "0(待首条)" + else + printf ' %-22s %8s\n' "$f" "无" + fi +done +for f in shadow_books.jsonl.gz shadow_tape.jsonl.gz; do + p="$OUT/$f" + [[ -f "$p" ]] && printf ' %-22s %8s\n' "$f" "$(du -h "$p" | cut -f1)" \ + || printf ' %-22s %8s\n' "$f" "无" +done + +say "各站点行数(确认 site 列生效)" +p="$OUT/shadow_latency.csv" +if [[ -f "$p" ]]; then + awk -F, 'NR>1{c[$1]++} END{for(s in c) printf " %-16s %8d 行\n", s, c[s]}' "$p" +else + echo " 尚无数据" +fi + +say "到达延迟中位(本站,跨地对比的主指标)" +# 按列名取下标,不写死数字:加了 site 列之后字段整体右移过一次, +# 写死 $8 会静默变成读 lag_signal_ms +if [[ -s "$p" ]]; then + for sym in BTC ETH SOL; do + med=$(awk -F, -v s="$sym" ' + NR==1 { for (i=1;i<=NF;i++) { if ($i=="sym") si=i; if ($i=="lag_data_ms") li=i } ; next } + $si==s && $li!="" { print $li }' "$p" | sort -n | awk ' + { v[NR]=$1 } END { if (NR) printf "%.0f %d", v[int((NR+1)/2)], NR }') + [[ -n "$med" ]] && printf ' %-5s 中位 %6sms (n=%s)\n' "$sym" ${med} \ + || printf ' %-5s 尚无数据\n' "$sym" + done + echo " 参考:本机(新加坡)补丁后实测 350~650ms,理论下限约 500ms" +fi diff --git a/research/live/shadow_hb.py b/research/live/shadow_hb.py index 45028eb..708ba86 100644 --- a/research/live/shadow_hb.py +++ b/research/live/shadow_hb.py @@ -58,6 +58,8 @@ import csv import gzip import json import math +import os +import socket import time from collections import deque from concurrent.futures import ProcessPoolExecutor @@ -69,6 +71,10 @@ import pandas as pd from lib.shadow_budget import LAG_ALARM_MS, LAG_WINDOW, lag_healthy +# 站点标识。跨地对比时两台机器的 CSV 要能合起来读,没有这一列就分不清哪行 +# 来自哪台。默认取主机名,部署脚本会显式传 SHADOW_SITE(如 sg-hetzner) +SITE = os.environ.get("SHADOW_SITE") or socket.gethostname() + SYMS = ("BTC", "ETH", "SOL") # 多存一根:deque 尾部是尚未收盘的当前根,剔除后正好剩 step39 定下的窗口 LTF_BARS, HTF_BARS = 2001, 801 # 有效窗口 2000 / 800,命中率在此饱和 @@ -108,12 +114,30 @@ def _writer(path: Path, cols: list[str]): flush=True) fresh = not path.exists() or path.stat().st_size == 0 f = path.open("a", newline="") - w = csv.DictWriter(f, fieldnames=cols) + w = _SiteWriter(csv.DictWriter(f, fieldnames=cols)) if fresh: w.writeheader() return f, w +class _SiteWriter: + """DictWriter 的薄包装,自动补上 site 列。 + + 逐个 writerow 手加 site 有四处,漏一处就是静默的空值,而跨地对比正是靠 + 这一列区分数据来源。在这里注入,漏不掉。 + """ + + def __init__(self, w: csv.DictWriter) -> None: + self._w = w + + def writeheader(self) -> None: + self._w.writeheader() + + def writerow(self, row: dict) -> None: + row.setdefault("site", SITE) + self._w.writerow(row) + + class BookLog: """把完整盘口快照落成 gzip JSONL。 @@ -136,7 +160,7 @@ class BookLog: asks: np.ndarray) -> None: # 只留价与量两列,update_id 对离线分析没用。round 到 10 位避免 # float repr 把文件撑大一倍 - rec = {"sym": sym, "kline_ts": kline_ts, "label": label, + rec = {"site": SITE, "sym": sym, "kline_ts": kline_ts, "label": label, "delay_ms": delay_ms, "target": target, "book_ts": book_ts, "bids": [[round(float(p), 10), round(float(a), 10)] for p, a, *_ in bids], @@ -187,7 +211,7 @@ class TapeLog: d = self.acc.get(sym) if not d or (not d["b"] and not d["s"]): return - rec = {"sym": sym, "bar_ts": bar_ts, + rec = {"site": SITE, "sym": sym, "bar_ts": bar_ts, "buys": {f"{p:.10g}": round(v, 10) for p, v in sorted(d["b"].items())}, "sells": {f"{p:.10g}": round(v, 10) @@ -291,7 +315,7 @@ class Shadow: d = out_dir() # 追加模式:长跑期间若重启,已收集的样本不该被清掉 self.f_sig, self.w_sig = _writer(d / "shadow_signals.csv", [ - "sym", "kline_ts", "direction", + "site", "sym", "kline_ts", "direction", "h1_agree", "ladder_ok", "gate_ok", "pass_all", "lag_ok", "atr_pct", "atr_bp", "t_close_ms", "t_data_ms", "t_signal_ms", @@ -301,12 +325,12 @@ class Shadow: "baseline_px", "mid", "best_px", "fill_px", "filled", "depth_ok", "slip_bp", "drift_bp", "spread_bp", "impact_bp"]) self.f_lat, self.w_lat = _writer(d / "shadow_latency.csv", [ - "sym", "kline_ts", "t_close_ms", "t_data_ms", "t_signal_ms", + "site", "sym", "kline_ts", "t_close_ms", "t_data_ms", "t_signal_ms", "lag_data_ms", "lag_signal_ms", "compute_ms", "n_bars", "n_hits", "n_pass", "atr_bp", "lag_med_ms", "lag_ok"]) # 无条件漂移:每根都记,用来和信号根上的条件漂移对照 self.f_drf, self.w_drf = _writer(d / "shadow_drift.csv", [ - "sym", "kline_ts", "delay_label", "delay_ms", + "site", "sym", "kline_ts", "delay_label", "delay_ms", "book_ts", "book_lag_ms", "baseline_px", "mid", "drift_bp_long"]) # 完整深度。挂在无条件漂移那条路径上,所以每根 K 线的四个固定延迟点 # 都有一份,信号根上再补一份 actual 点