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

Spark窗口函数处理含大量列的ML DataFrame前向填充问题

处理超大规模CSV的前向填充(ffill)方案

嘿,处理1.6亿行、500+特征的超大规模CSV做前向填充,直接硬怼Pandas全量加载绝对会爆内存,得用更聪明的办法。下面给你几个实战性强的方案,从单机到分布式都有:

方案一:Pandas分块处理(单机可用,内存压力可控)

如果你的机器内存还能扛住单块数据(比如100万行),可以用分块读取+跨块续接的方式实现前向填充,核心是保存上一块的最后一行有效值,用来填充下一块开头的空值:

import pandas as pd

# 分块大小根据你的内存调整,比如100万行/块
chunk_size = 1_000_000
# 保存上一块的最后有效值,用于跨块填充
last_valid_row = None
output_path = "filled_data.csv"
first_write = True

# 务必确保数据已经按Timestamp升序排序!否则填充逻辑完全错误
for chunk in pd.read_csv("your_raw_data.csv", chunksize=chunk_size):
    # 用前一块的最后有效值填充当前块开头的空值
    if last_valid_row is not None:
        for col in chunk.columns:
            if pd.isna(chunk.iloc[0][col]):
                chunk[col].fillna(last_valid_row[col], limit=1, inplace=True)
    
    # 对当前块执行前向填充
    chunk = chunk.ffill()
    
    # 更新上一块的最后有效值
    last_valid_row = chunk.iloc[-1].copy()
    
    # 写入结果,第一次写入带表头,之后追加不带
    chunk.to_csv(output_path, mode="a", header=first_write, index=False)
    first_write = False

注意点:

  • 必须先确认数据是按Timestamp升序排序的,如果没排序,先通过命令行工具(比如sort)或者分块排序的方式预处理,否则前向填充毫无意义
  • 分块大小别设太大,避免单块数据占满内存

方案二:Dask(单机多核/集群,语法和Pandas一致)

Dask是专门为大数据设计的框架,语法和Pandas几乎完全兼容,自动帮你处理分块和跨块逻辑,还能利用多核CPU加速:

import dask.dataframe as dd

# 读取CSV,Dask自动分块
ddf = dd.read_csv("your_raw_data.csv")

# 先按Timestamp排序(如果数据还没排序的话)
ddf = ddf.sort_values("Timestamp")

# 执行前向填充,和Pandas的ffill用法完全一样
ddf_filled = ddf.ffill()

# 保存结果,single_file=True会生成单个CSV,否则会生成多个分块文件
ddf_filled.to_csv("filled_dask_data.csv", single_file=True, index=False)

优势:

  • 不用手动处理跨块填充的细节,Dask内部自动搞定
  • 支持多核并行处理,速度比单进程Pandas快很多
  • 内存占用远低于Pandas全量加载

方案三:PySpark(分布式集群,超大规模数据)

如果单机完全扛不住,直接上Spark做分布式处理,适合TB级别的数据:

from pyspark.sql import SparkSession
from pyspark.sql.functions import last
from pyspark.sql.window import Window
import sys

# 初始化Spark会话
spark = SparkSession.builder.appName("FFillLargeData").getOrCreate()

# 读取CSV,注意指定header和inferSchema(或者手动指定Schema提高效率)
df = spark.read.csv("your_raw_data.csv", header=True, inferSchema=True)

# 定义窗口:如果是全局时间序列,直接按Timestamp排序;如果是按ID分组的时间序列,要加partitionBy("ID")
# rowsBetween(-sys.maxsize, 0)表示取当前行及之前所有行的最后一个非空值
window_spec = Window.orderBy("Timestamp").rowsBetween(-sys.maxsize, 0)
# 如果是分组场景,改成下面这行:
# window_spec = Window.partitionBy("ID").orderBy("Timestamp").rowsBetween(-sys.maxsize, 0)

# 对每个特征列执行前向填充
filled_df = df.select(
    "ID", "Timestamp",
    *[last(col, ignorenulls=True).over(window_spec).alias(col) for col in df.columns if col not in ["ID", "Timestamp"]]
)

# 保存结果,mode="overwrite"会覆盖已有文件
filled_df.write.csv("filled_spark_data.csv", header=True, mode="overwrite")

注意点:

  • 如果你的数据是按ID分组的(每个ID有独立的时间序列),一定要用partitionBy("ID"),不然会把不同ID的序列混在一起填充
  • 可以手动指定Schema来提高读取效率,避免Spark自动推断Schema的开销

关键前置提示

  1. 排序优先:所有方案的前提都是数据按Timestamp升序排列,否则前向填充的逻辑完全错误
  2. 空值预处理:先抽样检查各列的空值分布,对于全空的列,ffill也无法填充,建议单独标记或删除
  3. 小样本测试:先拿几百行带空值的测试数据验证代码逻辑,确认填充结果正确后再跑全量数据

内容的提问来源于stack exchange,提问作者cristi.calugaru

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:15:07