Python Polars八表连接内存耗尽问题求助
Polars八表连接内存耗尽问题排查与解决方案探讨
问题背景
这是一项POC测试,目标是验证Polars是否比现有SQL方案更快、更优、更经济。首个测试用例为八表连接的count(*)统计,八张表已加载为Lazy DataFrame,预期结果约3000万,但在配备NVIDIA L4 GPU的4核16GB系统中因内存耗尽被终止(该系统配置优于当前SQL测试平台)。
已尝试操作
- 多数表列数较少(<10列),部分表列数较多(>30列),已尝试仅保留连接所需列,但未解决内存问题
- 移除
collect(engine="gpu")选项后,进程不会终止,能正常返回执行计划 - 仅2-3张表连接时无异常,加入大表后出现内存耗尽
- 确认是系统16GB内存耗尽(GPU内存64GB未被占满),使用
.collect(streaming=True)、.collect(engine="gpu")、.collect()均出现相同问题 - 已尽可能压缩数据类型,但仍存在大量String类型字段
- 使用Polars 1.24版本,尝试
new_streaming=True时提示“not implemented”,使用streaming=True收到弃用警告
解决方案建议
- 优化连接顺序,缩减中间结果
- 优先对大表做预过滤,减少参与连接的数据量;调整连接顺序,先处理基数较低的关联键,避免中间结果过度膨胀。比如先完成
noun_coupon_nk相关的表连接,再处理noun_airport_nk的多表关联。
- 优先对大表做预过滤,减少参与连接的数据量;调整连接顺序,先处理基数较低的关联键,避免中间结果过度膨胀。比如先完成
- 压缩String类型内存占用
- 将重复率较高的String字段转换为
pl.Categorical类型,大幅降低内存消耗。示例代码:noun1 = noun1.with_columns(pl.col("noun_airport_nk").cast(pl.Categorical))
- 将重复率较高的String字段转换为
- 提前聚合,避免全量连接
- 不需要生成完整的连接结果再计数,可以通过统计各关联键的唯一值数量,结合关联关系计算最终count,跳过全表连接步骤。比如:
# 统计各表关联键的唯一值数量 noun1_unique = noun1.lazy().select("noun_airport_nk").unique().count() adj2_unique = adj2.lazy().select("noun_airport_nk").unique().count() # 根据关联逻辑计算最终count,避免生成全量连接数据
- 不需要生成完整的连接结果再计数,可以通过统计各关联键的唯一值数量,结合关联关系计算最终count,跳过全表连接步骤。比如:
- 升级Polars版本
- 升级到最新版Polars,新版本对内存管理和流式处理有优化,可能解决
new_streaming=True未实现的问题,同时修复旧版本的内存泄漏或低效问题。
- 升级到最新版Polars,新版本对内存管理和流式处理有优化,可能解决
- 系统层面临时调整
- 临时增加系统内存,验证是否为内存瓶颈导致的问题;也可尝试设置Polars内存限制,触发流式处理(需新版本支持):
pl.set_memory_limit("12GB")
- 临时增加系统内存,验证是否为内存瓶颈导致的问题;也可尝试设置Polars内存限制,触发流式处理(需新版本支持):
测试代码
lazy_final = ( noun1.lazy().select(["noun_airport_nk"]) .join(adj2.lazy().select(["noun_airport_nk"]), left_on="noun_airport_nk", right_on="noun_airport_nk", how="inner", suffix="__adj2") .join(verb3.lazy().select(["noun_airport_nk","noun_coupon_nk"]), on="noun_airport_nk", how="inner", suffix="__verb3") .join(noun4.lazy().select(["noun_coupon_nk"]), left_on="noun_coupon_nk", right_on="noun_coupon_nk", how="inner", suffix="__noun4") .join(adj5.lazy().select(["noun_coupon_nk"]), on="noun_coupon_nk", how="inner", suffix="__adj5") .join(verb6.lazy().select(["noun_airport_nk"]), on="noun_airport_nk", how="inner", suffix="__verb6") .join(noun7.lazy().select(["noun_airport_nk"]), on="noun_airport_nk", how="inner", suffix="__noun7") .join(adj8.lazy().select(["noun_airport_nk"]), on="noun_airport_nk", how="inner", suffix="__adj8") ) count_df = lazy_final.select(pl.len().alias('count')) print(count_df.collect(engine="gpu"))
内容的提问来源于stack exchange,提问作者sicsmpr
相关产品推荐
相关产品推荐

