
本文介绍一种纯 Polars 表达式驱动的高效方案,通过交叉连接 + 条件过滤 + 分组聚合,在大型 DataFrame 中快速计算每行到首个满足 B ≥ current_B × 1.5 的后续行的行距(以行索引差表示)。
本文介绍一种纯 polars 表达式驱动的高效方案,通过交叉连接 + 条件过滤 + 分组聚合,在大型 dataframe 中快速计算每行到首个满足 `b ≥ current_b × 1.5` 的后续行的行距(以行索引差表示)。
在处理大规模结构化数据时,避免 Python 循环、充分利用 Polars 的惰性执行与查询优化至关重要。本例目标是:对每一行,找出其后第一个满足 Column B ≥ 当前行 Column B × 1.5 的记录,并将二者在 Column A(作为逻辑行序标识)上的差值写入新列 Column C;若不存在则置为 null。
核心思路是将问题转化为「自连接 + 条件筛选 + 最近匹配」:
- 利用 cross 连接生成所有 (i, j) 行对(i 为当前行,j 为候选后续行);
- 筛选满足 A_j > A_i(确保是后续行)且 B_j ≥ 1.5 × B_i 的配对;
- 对每个 A_i 分组,取 A_j 最小者(即最近的满足条件行),计算行距 A_j − A_i。
以下是完整实现(推荐始终启用 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"),
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") # 取最小 A_right(即最近匹配)
)
.select(
"Column A",
(pl.col("next_A") - pl.col("Column A")).alias("Column C")
),
on="Column A",
how="left"
)
.collect(streaming=True)
)
print(result)✅ 关键优势说明:
- 零 Python 循环:全程基于 Polars 原生表达式,可被查询优化器重写、并行化;
- 内存友好:streaming=True 启用流式处理,避免全量交叉连接导致的 O(n²) 内存峰值;
- 语义清晰:逻辑直译业务需求(“找下一个满足条件的行”),易于维护与调试;
- 可扩展性强:支持任意数值阈值(如 1.3、2.0)或复杂条件组合(如 B_j ≥ B_i * 1.5 and C_j > 10)。
⚠️ 注意事项:
- Column A 必须严格递增且唯一(作为行序锚点),若原始索引不满足,请先 .with_row_index() 构建序号列;
- 若数据量极大(千万级+),建议配合 allow_parallel=True(默认启用)及合理 chunk size;
- cross 连接虽经优化,仍建议在 filter 前尽可能缩小右表(如添加时间窗口限制 A_right <= A + 100)进一步提速。
该方案兼顾正确性、性能与可读性,是 Polars 生态中解决“向前查找首个匹配”类问题的标准范式。

















