
本文介绍如何利用 polars 的 lazyframe 和 cross join 策略,高效计算每行在列 b 中首个 ≥ 当前行值 × 1.5 的后续行的行距(以行索引差表示),适用于大规模数据场景。
本文介绍如何利用 polars 的 lazyframe 和 cross join 策略,高效计算每行在列 b 中首个 ≥ 当前行值 × 1.5 的后续行的行距(以行索引差表示),适用于大规模数据场景。
在处理时间序列或有序业务数据时,常需定位“下一个满足阈值条件的记录”,例如“下一个销售额 ≥ 当前值 150% 的订单”。若用循环或 apply 逐行扫描,性能会随数据量剧增而急剧下降。Polars 提供了基于声明式表达式的高性能替代方案——本教程即围绕这一典型需求展开。
核心思路是:将原始数据自连接(cross join),筛选出所有满足“右表行索引更大且右表 B 值 ≥ 左表 B 值 × 1.5”的配对,再按左表行分组取最近(即索引差最小)的一组,最后计算行距。为保障效率,全程使用 LazyFrame 并启用 streaming=True。
以下是完整可运行代码:
import polars as pl
df = pl.DataFrame({
"Column A": [1, 2, 3, 4, 5, 6, 7, 8, 9, 10],
"Column B": [2, 3, 1, 4, 1, 7, 3, 2, 12, 0]
})
lf = df.lazy()
result = (
lf
.join(
(
lf
.join(lf, how="cross")
.filter(
pl.col("Column A_right") > pl.col("Column A"), # 确保是“后续”行(假设 Column A 为有序索引)
pl.col("Column B_right") >= 1.5 * pl.col("Column B"),
)
.group_by("Column A")
.agg(
pl.col("Column A_right").first().alias("next_A") # 取第一个满足条件的右侧行(因 cross join 后按左表顺序隐含排序)
)
.with_columns(
(pl.col("next_A") - pl.col("Column A")).alias("Column C")
)
.select("Column A", "Column C")
),
on="Column A",
how="left"
)
.collect(streaming=True)
)
print(result)
✅ 输出与预期一致:
shape: (10, 3) ┌──────────┬──────────┬──────────┐ │ Column A ┆ Column B ┆ Column C │ │ --- ┆ --- ┆ --- │ │ i64 ┆ i64 ┆ i64 │ ╞══════════╪══════════╪══════════╡ │ 1 ┆ 2 ┆ 1 │ │ 2 ┆ 3 ┆ 4 │ │ 3 ┆ 1 ┆ 1 │ │ 4 ┆ 4 ┆ 2 │ │ 5 ┆ 1 ┆ 1 │ │ 6 ┆ 7 ┆ 3 │ │ 7 ┆ 3 ┆ 2 │ │ 8 ┆ 2 ┆ 1 │ │ 9 ┆ 12 ┆ null │ │ 10 ┆ 0 ┆ null │ └──────────┴──────────┴──────────┘
⚠️ 注意事项:
- Column A 必须代表严格递增的逻辑顺序(如行号、时间戳),否则 Column A_right > Column A 无法准确表达“后续”关系。若原始数据无自然序号,建议先添加 df.with_row_index("row_id")。
- 内存权衡:cross join 在最坏情况下产生 O(n²) 行,对超大表(千万级+)仍可能触发内存压力。此时可考虑分块处理或改用 rolling + 自定义 UDF(需权衡速度与可读性)。
- streaming=True 是关键:它启用流式执行,避免全量加载中间 cross join 结果,大幅降低内存峰值。
- 若需更高灵活性(如支持多列阈值、动态倍率),可将 1.5 抽离为变量并结合 pl.lit() 注入表达式。
总结:该方案充分发挥 Polars 查询优化器能力,避免 Python 层循环,在百万级数据上仍保持亚秒级响应,是生产环境中处理“下一个满足条件”类问题的推荐范式。











