修掉 gzip 追加会毁掉整个文件的数据丢失,并分开买卖两侧的成交分布曲线

两件事,都是「静默出错」那一类。

一、BookLog/TapeLog 追加到同一个 .gz,进程被 SIGKILL 时当前成员停在 deflate
块中间,下一轮追加的新成员接在垃圾字节之后。顺序解压在损坏点抛 invalid
block type,该点之后全部读不出来——包括后续每轮写进去的。而读侧的异常处理
把这个当成「正常的尾部截断」静默跳过,于是只读出 21 行还不报错。
原 docstring 里写的「只丢最后一个缓冲块,不会毁掉整个文件」是错的,已证伪。

写侧改成每轮运行一个文件;读侧按 gzip 成员边界扫描、坏成员单独跳过并出声
报告,同时把同前缀的多轮文件一并读入。旧损坏文件因此多恢复出 31/21 条
(tape)与 132/95 条(books)。

二、tape_shape 只统计主动买、只自区间顶部累积,这条曲线只适用于多头止盈。
exit_fill 两侧共用它,等于把空头的可成交量按多头分布高估。实测二者不对称:
主动买在顶部 20% 内已占 40%,主动卖在底部 20% 内只有 18%。分成 SHAPE_F 与
SHAPE_F_SHORT,avail_at 按方向查各自曲线。

两条曲线只有 45 根成交流样本,所以补了 --sensitivity:把空头可成交量砍一半,
10 万仓位下 BTC/ETH 预算完全不动、SOL 动 0.07bp。结论不依赖这 45 根样本。
顺带撤掉 step43 docstring 里已作废的「成交率 30%/16%/1.5%」。

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
jack
2026-08-28 03:23:07 +08:00
co-authored by Cursor
parent 54792fe015
commit 75fbf4167b
4 changed files with 183 additions and 37 deletions
+26 -5
View File
@@ -45,20 +45,37 @@ from lib.exit_model import FEE_MAKER, FEE_TAKER, SLIP
# 成交流实测的区间内成交分布形状(BTC/ETH/SOL 均值,见 shadow_depth.tape_shape # 成交流实测的区间内成交分布形状(BTC/ETH/SOL 均值,见 shadow_depth.tape_shape
SHAPE_K = np.linspace(0.0, 1.0, 21) SHAPE_K = np.linspace(0.0, 1.0, 21)
# 多头止盈挂卖出,靠**主动买**打上来 —— 自区间顶部向下累积
SHAPE_F = np.array([0.044, 0.088, 0.110, 0.179, 0.204, 0.240, 0.282, 0.316, SHAPE_F = np.array([0.044, 0.088, 0.110, 0.179, 0.204, 0.240, 0.282, 0.316,
0.390, 0.430, 0.465, 0.537, 0.583, 0.617, 0.662, 0.714, 0.390, 0.430, 0.465, 0.537, 0.583, 0.617, 0.662, 0.714,
0.761, 0.804, 0.857, 0.904, 1.000]) 0.761, 0.804, 0.857, 0.904, 1.000])
# 空头止盈挂买回,靠**主动卖**打下来 —— 自区间底部向上累积。
# 两侧并不对称:主动买集中在区间顶部(k=0.2 处已 40%),主动卖在底部只有
# 18%BTC/ETH 平均偏差 0.21。原先两侧共用 SHAPE_F,等于把空头的可成交量
# 按多头的分布高估。
# ⚠ 这条曲线只有 45 根样本(每币 14~16),够说明「不对称」这个方向,不够
# 定具体数值——一周数据到手要重新导出。用它而非 SHAPE_F 的理由是方向正确
# 优于数值精确:结论对形状不敏感(见 step43_fill_aware_budget 的敏感性检查),
# 但用错方向是系统性偏乐观。
SHAPE_F_SHORT = np.array([0.031, 0.076, 0.121, 0.193, 0.230, 0.265, 0.297,
0.356, 0.382, 0.420, 0.453, 0.485, 0.523, 0.552,
0.586, 0.629, 0.681, 0.754, 0.813, 0.871, 1.000])
TAKER_SHARE = 0.5 TAKER_SHARE = 0.5
def avail_at(price: float, hi: float, lo: float, vol_notional: float, def avail_at(price: float, hi: float, lo: float, vol_notional: float,
is_long: bool, kgrid=SHAPE_K, f=SHAPE_F) -> float: is_long: bool, kgrid=SHAPE_K, f=None,
f_short=None) -> float:
"""该根里能打到限价 `price` 的对手方成交额。 """该根里能打到限价 `price` 的对手方成交额。
多头在 `price` 挂卖出,靠价格 ≥ price 的主动买成交;空头挂买回,靠 多头在 `price` 挂卖出,靠价格 ≥ price 的主动买成交;空头挂买回,靠
价格 ≤ price 的主动卖成交。方向用显式参数而非「把价格取负」——取负会 价格 ≤ price 的主动卖成交。方向用显式参数而非「把价格取负」——取负会
让所有价格变成负数,任何对价格正负的假设都会静默失效。 让所有价格变成负数,任何对价格正负的假设都会静默失效。
两个方向查各自的成交分布曲线(实测二者不对称,见 SHAPE_F_SHORT)。
""" """
f = SHAPE_F if f is None else f
f_short = SHAPE_F_SHORT if f_short is None else f_short
if vol_notional <= 0 or not (np.isfinite(hi) and np.isfinite(lo)): if vol_notional <= 0 or not (np.isfinite(hi) and np.isfinite(lo)):
return 0.0 return 0.0
if hi <= lo: if hi <= lo:
@@ -77,16 +94,20 @@ def avail_at(price: float, hi: float, lo: float, vol_notional: float,
if price >= hi: if price >= hi:
return vol_notional # 整根都在限价之下 return vol_notional # 整根都在限价之下
k = (price - lo) / (hi - lo) k = (price - lo) / (hi - lo)
return float(np.interp(k, kgrid, f)) * vol_notional return float(np.interp(k, kgrid, f if is_long else f_short)) * vol_notional
def walk_filled(cdf: pd.DataFrame, sig: pd.DataFrame, notional: float, def walk_filled(cdf: pd.DataFrame, sig: pd.DataFrame, notional: float,
sl: float = 2.0, scale_at: float = 3.0, runner: float = 8.0, sl: float = 2.0, scale_at: float = 3.0, runner: float = 8.0,
runner_stop: float = 2.0, maxb: int = 48, runner_stop: float = 2.0, maxb: int = 48,
taker_share: float = TAKER_SHARE) -> pd.DataFrame: taker_share: float = TAKER_SHARE,
f=None, f_short=None) -> pd.DataFrame:
"""前推每笔信号,返回按成交量结算的出场权重与毛收益。 """前推每笔信号,返回按成交量结算的出场权重与毛收益。
每行的 `w_*` 是各出场去向占**全仓名义额**的比例,四者相加为 1。 每行的 `w_*` 是各出场去向占**全仓名义额**的比例,四者相加为 1。
`f` / `f_short` 是两个方向的区间内成交分布曲线,留出接口是为了能做敏感性
检查——这两条曲线目前只有 45 根样本,必须能验证结论对它们不敏感。
""" """
high = cdf["high"].to_numpy(float) high = cdf["high"].to_numpy(float)
low = cdf["low"].to_numpy(float) low = cdf["low"].to_numpy(float)
@@ -130,7 +151,7 @@ def walk_filled(cdf: pd.DataFrame, sig: pd.DataFrame, notional: float,
v = vol[j] * close[j] * taker_share v = vol[j] * close[j] * taker_share
is_long = d == 1 is_long = d == 1
if rem_scale > 0: if rem_scale > 0:
got = avail_at(p_scale, hi, lo, v, is_long) got = avail_at(p_scale, hi, lo, v, is_long, SHAPE_K, f, f_short)
fill = min(rem_scale, got / notional) if notional > 0 else \ fill = min(rem_scale, got / notional) if notional > 0 else \
rem_scale rem_scale
if fill > 0: if fill > 0:
@@ -138,7 +159,7 @@ def walk_filled(cdf: pd.DataFrame, sig: pd.DataFrame, notional: float,
w_scale += fill w_scale += fill
scaled_any = True scaled_any = True
if rem_run > 0: if rem_run > 0:
got = avail_at(p_run, hi, lo, v, is_long) got = avail_at(p_run, hi, lo, v, is_long, SHAPE_K, f, f_short)
fill = min(rem_run, got / notional) if notional > 0 else \ fill = min(rem_run, got / notional) if notional > 0 else \
rem_run rem_run
if fill > 0: if fill > 0:
+72 -21
View File
@@ -28,8 +28,8 @@
from __future__ import annotations from __future__ import annotations
import argparse import argparse
import gzip
import json import json
import zlib
from pathlib import Path from pathlib import Path
import numpy as np import numpy as np
@@ -41,27 +41,69 @@ def out_dir() -> Path:
return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out" return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out"
def read_jsonl_gz(path: Path): GZ_MAGIC = b"\x1f\x8b\x08"
"""逐行读 gzip JSONL,末尾截断则静默停止。
采集仍在进行时,最后一个 gzip 成员缺结尾标记;不接这个异常的话,
整个分析会因为文件尾而失败,前面几万条完好记录一起丢掉。 def _members(blob: bytes):
"""把可能损坏的多成员 gzip 拆成「逐个成员解压」,坏成员跳过。
采集进程被 SIGKILL 时,当前 gzip 成员停在 deflate 块中间、没有结尾标记。
下一次运行以追加方式写入的新成员就接在这段垃圾字节后面。此时用
`gzip.open` 顺序读会在损坏点抛 `invalid block type`**该点之后的所有
数据都读不出来**——包括后续每一轮运行写进去的。曾因此只读出 21 行而
误以为样本就那么少,且不报错。
所以按成员边界扫描:某个成员解压失败,只丢它,然后前进到下一个 magic
继续。返回 (解压出的字节, 跳过的成员数)。
""" """
if not path.exists(): out, skipped, i, n = [], 0, blob.find(GZ_MAGIC), len(blob)
while 0 <= i < n:
d = zlib.decompressobj(16 + zlib.MAX_WBITS)
try:
chunk = d.decompress(blob[i:])
except zlib.error:
chunk = b""
if chunk:
out.append(chunk)
# 成员完整时 unused_data 指向下一成员;否则只能往前找 magic
if d.eof and d.unused_data:
nxt = n - len(d.unused_data)
else:
nxt = blob.find(GZ_MAGIC, i + 3)
if chunk == b"":
skipped += 1
i = nxt if nxt > i else -1
return b"".join(out), skipped
def read_jsonl_gz(path: Path, quiet: bool = False):
"""读 gzip JSONL,跨运行文件汇总,坏成员跳过而非静默截断。
`path` 既可以是单个文件,也当作前缀用:同目录下 `<stem>.*.jsonl.gz`
(每轮运行一个)会一并读入,这样重启不再把历史数据连坐。
"""
base = path.name.replace(".jsonl.gz", "")
files = sorted({*path.parent.glob(f"{base}.*.jsonl.gz"),
*([path] if path.exists() else [])})
if not files:
return return
n_ok = 0 for f in files:
try: blob = f.read_bytes()
with gzip.open(path, "rt", encoding="utf-8") as fh: data, skipped = _members(blob)
for line in fh: n_ok = n_bad = 0
try: for line in data.split(b"\n"):
rec = json.loads(line) if not line:
except json.JSONDecodeError: continue
break # 半行,说明写到这里被打断 try:
n_ok += 1 rec = json.loads(line)
yield rec except (json.JSONDecodeError, UnicodeDecodeError):
except (EOFError, OSError, gzip.BadGzipFile): n_bad += 1 # 损坏边界上的半行
# 采集进程正在写,尾部不完整属正常 continue
pass n_ok += 1
yield rec
if not quiet and (skipped or n_bad):
print(f"{f.name}: 读出 {n_ok:,} 条,"
f"跳过 {skipped} 个损坏成员 / {n_bad} 个半行")
def impact_bp(levels: list, notional: float, mid: float) -> float | None: def impact_bp(levels: list, notional: float, mid: float) -> float | None:
@@ -239,8 +281,17 @@ def maker_fill(tape_path: Path, mults=(3.0, 8.0),
def tape_shape(tape_path: Path, kgrid: np.ndarray) -> dict[str, np.ndarray]: def tape_shape(tape_path: Path, kgrid: np.ndarray) -> dict[str, np.ndarray]:
"""成交流给「形状」:一根的主动买成交额里,有多少比例落在区间顶部 k 之内。 """成交流给「形状」:一根的主动买成交额里,有多少比例落在区间顶部 k 之内。
形状与规模分开是为了绕开成交流样本小的限制——形状是微观结构性质, 形状与规模分开是为了绕开成交流样本小的限制——规模(每根成交多少钱)由
几十根就相当稳定;规模(每根成交多少钱)则由 210 天历史成交量提供。 210 天历史成交量提供,成交流只需给出形状
⚠ 只算主动**买**、只自顶部累积,所以这条曲线只适用于**多头**止盈(挂卖
出,靠主动买打上来)。空头止盈要用主动卖自底部累积的曲线,二者实测并不
对称:主动买在区间顶部 20% 内已占 40%,主动卖在底部 20% 内只有 18%
BTC/ETH 平均偏差 0.21。lib/exit_fill 里两条曲线是分开的(SHAPE_F 与
SHAPE_F_SHORT);用同一条会把空头的可成交量按多头分布高估。
「形状几十根就稳定」这个说法要打折:45 根样本足以看出上面那个方向性差异,
但不足以定数值。好在预算对形状不敏感(见 step43 --sensitivity)。
""" """
acc: dict[str, list[np.ndarray]] = {} acc: dict[str, list[np.ndarray]] = {}
for r in read_jsonl_gz(tape_path): for r in read_jsonl_gz(tape_path):
+25 -5
View File
@@ -96,6 +96,20 @@ def out_dir() -> Path:
return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out" return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out"
RUN_ID = time.strftime("%Y%m%dT%H%M%S", time.gmtime())
def run_path(path: Path) -> Path:
"""把 `x.jsonl.gz` 变成本轮专属的 `x.<RUN_ID>.jsonl.gz`。
gzip 追加流在进程被杀后不可靠(见 BookLog 的说明),分文件是唯一能保证
历史数据不被后续运行连坐的办法。读侧 shadow_depth.read_jsonl_gz 会把
同前缀的所有文件一并读入,所以分文件对分析是透明的。
"""
return path.with_name(path.name.replace(".jsonl.gz",
f".{RUN_ID}.jsonl.gz"))
def _writer(path: Path, cols: list[str]): def _writer(path: Path, cols: list[str]):
"""追加模式打开;表头对不上就先把旧文件归档。 """追加模式打开;表头对不上就先把旧文件归档。
@@ -146,13 +160,16 @@ class BookLog:
一次采集回答所有资金量级的问题——包括容量上限那个必须现在就算、 一次采集回答所有资金量级的问题——包括容量上限那个必须现在就算、
不该等实盘暴露的数。 不该等实盘暴露的数。
用 gzip 追加(多个 gzip 成员首尾相接仍可正常解压),进程被杀也只丢最后 **每轮运行单独一个文件**,不追加到同一个。追加看着更省事,实际很危险:
一个缓冲块,不会毁掉整个文件。 进程被 SIGKILL 时当前 gzip 成员停在 deflate 块中间,下一轮追加的新成员
接在这段垃圾字节之后,顺序解压会在损坏点抛 `invalid block type`,该点
之后的所有数据——包括后续每一轮写进去的——全都读不出来。已经因此丢过
一次。分文件后损坏最多只影响被杀那一轮的尾部。
""" """
def __init__(self, path: Path) -> None: def __init__(self, path: Path) -> None:
self.path = path self.path = run_path(path)
self.fh = gzip.open(path, "at", encoding="utf-8") self.fh = gzip.open(self.path, "at", encoding="utf-8")
self.n = 0 self.n = 0
def write(self, sym: str, kline_ts: int, label: str, delay_ms: int, def write(self, sym: str, kline_ts: int, label: str, delay_ms: int,
@@ -192,10 +209,13 @@ class TapeLog:
聚合到「根 × 价位」而不是逐笔:判据是「本根内有多少量在 ≥ 限价处成交」, 聚合到「根 × 价位」而不是逐笔:判据是「本根内有多少量在 ≥ 限价处成交」,
逐笔的时序对这个判据没有增量信息,而聚合能把体量压下两个数量级。 逐笔的时序对这个判据没有增量信息,而聚合能把体量压下两个数量级。
每轮运行单独一个文件,理由同 BookLog。
""" """
def __init__(self, path: Path) -> None: def __init__(self, path: Path) -> None:
self.fh = gzip.open(path, "at", encoding="utf-8") self.path = run_path(path)
self.fh = gzip.open(self.path, "at", encoding="utf-8")
# sym -> side('b'/'s') -> price -> 累计基础币量 # sym -> side('b'/'s') -> price -> 累计基础币量
self.acc: dict[str, dict[str, dict[float, float]]] = {} self.acc: dict[str, dict[str, dict[float, float]]] = {}
self.n_trades = 0 self.n_trades = 0
+60 -6
View File
@@ -1,17 +1,29 @@
"""Step 43:把「限价单全额成交」的假设换成按成交量结算,重算滑点预算。 """Step 43:把「限价单全额成交」的假设换成按成交量结算,重算滑点预算。
step42 的预算BTC 8.58 / ETH 20.64 / SOL 16.83bp建立在一个假设上:挂在 step42 的预算建立在一个假设上:挂在 3ATR 与 8ATR 的止盈限价单全额成交在目标
3ATR 与 8ATR 的止盈限价单全额成交在目标价。影子交易的成交流数据推翻了它—— 价。本脚本把它换成按成交流实测的成交量结算,让预算变成**仓位规模的函数**。
真实仓位下全额成交率只有 30%/16%/1.5%32 万仓位)。
预算因此不再是一个常数,而是**仓位规模的函数**。规模越大,止盈越难成交, ## 结论(2026-08-28
越多仓位被拖到止损或超时(taker,且吃滑点),预算越低。这条曲线与「冲击
反推的容量」是两个不同的约束,而后者宽松得多(100~500 万 vs 数万)。 **主口径 10 万 USDT 下成交率不是绑定约束。** 预算相对全额成交假设的降幅:
BTC 0% / ETH 0% / SOL 1.1%;即便到 100 万也只有 2.6% / 0.9% / 8.0%
⚠ 此前一版本文档写「真实仓位下全额成交率只有 30%/16%/1.5%」,那是
shadow_depth.composite_fill 的口径——只算首次触及那一根的可成交量,而真实
挂单在那儿常驻最多 48 根、每根都在成交。该数系统性偏悲观,已撤回。
另一个约束是冲击反推的容量上限(100~500 万),比成交率宽松得多。两者都不
绑定,所以 10 万仓位上限制来自别处,不是流动性。
结论对成交分布曲线**不敏感**:把空头侧可成交量砍一半,10 万仓位下 BTC/ETH
预算完全不动、SOL 动 0.07bp。见 `--sensitivity`。这一点重要,因为那两条曲线
目前只有 45 根成交流样本。
数据用 Bitget 210 天 1m,与影子测量同源同交易所。全量 366 万根峰值 24.5GB, 数据用 Bitget 210 天 1m,与影子测量同源同交易所。全量 366 万根峰值 24.5GB,
本机 15GB 跑不动;210 天 30 万根峰值约 2GB。 本机 15GB 跑不动;210 天 30 万根峰值约 2GB。
python research/step43_fill_aware_budget.py --syms BTC,ETH,SOL python research/step43_fill_aware_budget.py --syms BTC,ETH,SOL
python research/step43_fill_aware_budget.py --sensitivity
""" """
from __future__ import annotations from __future__ import annotations
@@ -92,13 +104,55 @@ def signals_for(sym: str, cache: Path) -> tuple[pd.DataFrame, pd.DataFrame]:
return cdf, sig return cdf, sig
def sensitivity(syms: list[str], cache: Path, notional: float = 1e5) -> None:
"""预算对成交分布曲线的敏感性。
SHAPE_F / SHAPE_F_SHORT 只有 45 根成交流样本,数值精度很低。所以必须先
证明结论对它们不敏感,否则整条预算曲线都建立在 45 根样本上。
三档:两侧共用买盘曲线(旧口径,空头偏乐观)/实测的方向各异曲线/把
空头可成交量再砍一半的悲观上界。
"""
from lib.exit_fill import SHAPE_F, SHAPE_F_SHORT, budget_bp, walk_filled
half = np.clip(SHAPE_F_SHORT * 0.5, 0.0, 1.0)
half[-1] = 1.0
cases = [("对称(旧口径)", SHAPE_F, SHAPE_F),
("实测不对称", SHAPE_F, SHAPE_F_SHORT),
("悲观:空头量减半", SHAPE_F, half)]
print(f"\n{'=' * 74}\n成交分布曲线敏感性 · 仓位 {notional:,.0f} USDT\n")
print(f" {'':<5}" + "".join(f"{n:>18}" for n, _, _ in cases))
for sym in syms:
try:
cdf, sig = signals_for(sym, cache)
except Exception as e:
print(f" {sym}: 跳过 {e!r}")
continue
row = f" {sym:<5}"
for _, fl, fs in cases:
r = walk_filled(cdf, sig, notional, SL, SCALE_AT, RUNNER,
RUNNER_STOP, MAXB, f=fl, f_short=fs)
row += f"{budget_bp(r):>13.2f}bp " if not r.empty \
else f"{'':>18}"
print(row)
del cdf
print("\n 2026-08-28 实测三档差异 ≤ 0.07bp,结论对曲线不敏感。")
def main() -> None: def main() -> None:
ap = argparse.ArgumentParser() ap = argparse.ArgumentParser()
ap.add_argument("--syms", default="BTC,ETH,SOL") ap.add_argument("--syms", default="BTC,ETH,SOL")
ap.add_argument("--cache", default="research/live/cache") ap.add_argument("--cache", default="research/live/cache")
ap.add_argument("--save", default="research/out/step43_fill_budget.csv") ap.add_argument("--save", default="research/out/step43_fill_budget.csv")
ap.add_argument("--sensitivity", action="store_true",
help="只跑成交分布曲线的敏感性检查")
a = ap.parse_args() a = ap.parse_args()
if a.sensitivity:
sensitivity(a.syms.split(","), Path(a.cache))
return
from lib.exit_fill import (assert_converges, budget_bp, net_bp, from lib.exit_fill import (assert_converges, budget_bp, net_bp,
walk_filled) walk_filled)
from lib.exit_model import cfg_name, slip_budget from lib.exit_model import cfg_name, slip_budget