You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Polars从超大规模时序数据集高效批量提取多段切片?

批量提取大CSV时序数据切片的Polars优化方案

问题描述

我需要从一个约25GB、8亿行的时序数据集中提取多段小切片。目前的实现代码如下:

from polars import pl

sample = pl.scan_csv(FILENAME, new_columns=["time", "force"]).slice(660_000_000, 3000).collect()

这段代码的执行耗时在0-5分钟不等,具体取决于切片的位置。若提取5段切片,总耗时约15分钟。由于Polars读取CSV时会扫描全量数据,我希望找到一种仅读取一次CSV就能批量提取所有所需切片的方法。链式调用多个slice显然不可行,是否有其他解决方案?

解决方案

方法1:行号过滤法(推荐)

利用Polars的with_row_index添加行号列,通过一次过滤筛选出所有目标切片的行,仅扫描CSV一次:

from polars import pl

# 定义需要提取的切片范围:(起始行号, 行数)
slice_ranges = [
    (660_000_000, 3000),
    (120_000_000, 2000),
    (780_000_000, 4000)
]

# 构建行号过滤表达式
filter_expr = None
for start, length in slice_ranges:
    end = start + length
    # 匹配[start, end-1]区间的行号
    range_expr = pl.col("row_nr").is_between(start, end - 1)
    filter_expr = range_expr if filter_expr is None else filter_expr | range_expr

# 一次性提取所有目标切片
result = (
    pl.scan_csv(FILENAME, new_columns=["time", "force"])
    .with_row_index("row_nr", offset=0)  # 行号从0开始,与slice索引对齐
    .filter(filter_expr)
    .drop("row_nr")
    .collect()
)

方法2:分块读取法(低内存场景)

如果内存有限,可手动分块读取CSV,仅保留与目标切片重叠的块内容:

from polars import pl

slice_ranges = [
    (660_000_000, 3000),
    (120_000_000, 2000),
    (780_000_000, 4000)
]

# 转换为左闭右开的区间
target_intervals = [(start, start + length) for start, length in slice_ranges]

chunk_size = 1_000_000  # 单次读取行数,可根据内存调整
current_row = 0
collected_slices = []

# 流式读取CSV分块
with pl.scan_csv(FILENAME, new_columns=["time", "force"]).streaming() as stream_reader:
    while True:
        chunk = stream_reader.slice(current_row, chunk_size).collect()
        if chunk.is_empty():
            break
        
        chunk_end = current_row + len(chunk)
        # 检查当前块与目标区间的重叠部分
        for start, end in target_intervals:
            overlap_start = max(start, current_row) - current_row
            overlap_end = min(end, chunk_end) - current_row
            if overlap_start < overlap_end:
                collected_slices.append(chunk.slice(overlap_start, overlap_end - overlap_start))
        
        current_row = chunk_end

# 合并所有提取的切片
result = pl.concat(collected_slices)

注意事项

  • 行号偏移:slice的起始索引从0开始,因此with_row_index需设置offset=0保证索引对应。
  • 性能对比:方法1依赖Polars的查询优化,代码简洁且效率高;方法2适合内存紧张的场景,避免一次性加载大量数据。

内容的提问来源于stack exchange,提问作者Jan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 02:59:59