fix(factor): load内存根治v2——类型转换下沉chunk级(dt解析/vt_symbol常量列/列裁剪),字符串三列活不过2100行小块,concat仅碰精简类型列;前轮lazy化未压住的6.6GB峰值源头即Utf8多份拷贝 [vps]

This commit is contained in:
2026-08-26 05:27:28 +08:00
parent 6eb977a0d0
commit 62673ea329
+17 -19
View File
@@ -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