From 62673ea329931a2d30c32670565b2c6cbb254639 Mon Sep 17 00:00:00 2001 From: claude_dev Date: Wed, 26 Aug 2026 05:27:28 +0800 Subject: [PATCH] =?UTF-8?q?fix(factor):=20load=E5=86=85=E5=AD=98=E6=A0=B9?= =?UTF-8?q?=E6=B2=BBv2=E2=80=94=E2=80=94=E7=B1=BB=E5=9E=8B=E8=BD=AC?= =?UTF-8?q?=E6=8D=A2=E4=B8=8B=E6=B2=89chunk=E7=BA=A7(dt=E8=A7=A3=E6=9E=90/?= =?UTF-8?q?vt=5Fsymbol=E5=B8=B8=E9=87=8F=E5=88=97/=E5=88=97=E8=A3=81?= =?UTF-8?q?=E5=89=AA),=E5=AD=97=E7=AC=A6=E4=B8=B2=E4=B8=89=E5=88=97?= =?UTF-8?q?=E6=B4=BB=E4=B8=8D=E8=BF=872100=E8=A1=8C=E5=B0=8F=E5=9D=97,conc?= =?UTF-8?q?at=E4=BB=85=E7=A2=B0=E7=B2=BE=E7=AE=80=E7=B1=BB=E5=9E=8B?= =?UTF-8?q?=E5=88=97;=E5=89=8D=E8=BD=AElazy=E5=8C=96=E6=9C=AA=E5=8E=8B?= =?UTF-8?q?=E4=BD=8F=E7=9A=846.6GB=E5=B3=B0=E5=80=BC=E6=BA=90=E5=A4=B4?= =?UTF-8?q?=E5=8D=B3Utf8=E5=A4=9A=E4=BB=BD=E6=8B=B7=E8=B4=9D=20[vps]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- sanguo_factor/universe.py | 36 +++++++++++++++++------------------- 1 file changed, 17 insertions(+), 19 deletions(-) diff --git a/sanguo_factor/universe.py b/sanguo_factor/universe.py index 2310d61..8ebd309 100644 --- a/sanguo_factor/universe.py +++ b/sanguo_factor/universe.py @@ -73,33 +73,31 @@ def load_universe_bars( finally: conn.close() if rows: - chunks.append(pl.DataFrame( - rows, - schema={"symbol": pl.Utf8, "exchange": pl.Utf8, "dt": pl.Utf8, - "volume": pl.Float64, "turnover": pl.Float64, "open_interest": pl.Float64, - "open": pl.Float64, "high": pl.Float64, "low": pl.Float64, "close": pl.Float64}, - orient="row", - )) + chunks.append( + pl.DataFrame(rows, schema={ + "symbol": pl.Utf8, "exchange": pl.Utf8, "dt": pl.Utf8, + "volume": pl.Float64, "turnover": pl.Float64, "open_interest": pl.Float64, + "open": pl.Float64, "high": pl.Float64, "low": pl.Float64, "close": pl.Float64, + }, orient="row") + .with_columns( + pl.col("dt").str.slice(0, 10).str.to_datetime("%Y-%m-%d").alias("datetime"), + pl.lit(f"{sym}.{ex}").alias("vt_symbol"), + pl.when(pl.col("volume") > 0) + .then(pl.col("turnover") / pl.col("volume")) + .otherwise(None).alias("vwap"), + ) + .select(["vt_symbol", "datetime", "open", "high", "low", "close", + "volume", "turnover", "open_interest", "vwap"]) + ) if not chunks: return _empty_alpha_df() df = pl.concat(chunks) - chunks.clear() # 立即释放逐symbol块(~1G),勿持有到函数尾 - # transform 链 lazy 单次物化:eager 逐步链每步全量拷贝,峰值3-4倍;lazy 融合后~1.5倍 + chunks.clear() df = ( df.lazy() - .with_columns( - pl.col("dt").str.slice(0, 10).str.to_datetime("%Y-%m-%d").alias("datetime"), - (pl.col("symbol") + "." + pl.col("exchange")).alias("vt_symbol"), - pl.when(pl.col("volume") > 0) - .then(pl.col("turnover") / pl.col("volume")) - .otherwise(None) - .alias("vwap"), - ) .sort(["vt_symbol", "datetime"]) .with_columns(pl.int_range(pl.len()).over("vt_symbol").alias("bar_idx")) - .select(["vt_symbol", "datetime", "open", "high", "low", "close", - "volume", "turnover", "open_interest", "vwap", "bar_idx"]) .collect() ) return df