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

如何将Dask Series转换为Dask Array并处理超大Parquet数据

处理大体积Snappy Parquet文件:Dask转Array并行预处理与TensorFlow分块训练

问题背景

处理超大规模.snappy.parquet文件,已通过Dask DataFrame读取并设置每个分区10行,得到的Dask Series中每个元素为2×125×125的numpy张量。直接compute()耗时极长,Pandas读取触发OOM,需转换为Dask Array实现并行处理,同时完成预处理并基于TensorFlow ResNet分块训练。


1. Dask Series转Dask Array

由于每个元素是固定形状的张量,可通过dask.array.from_delayed结合分区转换实现:

import dask.array as da
import dask.dataframe as dd
import numpy as np

# 读取数据(保留原分区配置)
train_path = ["/train.snappy.parquet"]
data = dd.read_parquet(
    train_path,
    engine="pyarrow",
    compression="snappy",
    columns=["X_jets", "y"],
    split_row_groups=10
)

# 定义单个分区的Series转数组函数
def series_to_array(series):
    # 堆叠分区内所有元素,得到形状为(10, 2, 125, 125)的数组
    return np.stack(series.values)

# 将每个分区转为延迟对象
delayed_partitions = data["X_jets"].map_partitions(
    series_to_array,
    meta=np.ndarray
).to_delayed()

# 拼接延迟对象为Dask Array
X_array = da.concatenate([
    da.from_delayed(
        d,
        shape=(10, 2, 125, 125),
        dtype=np.float32
    ) for d in delayed_partitions
])

# 同时将标签转为Dask Array
y_array = data["y"].to_dask_array(lengths=True)

2. 并行预处理优化

原预处理代码存在分区单独拟合StandardScaler导致统计量偏差的问题,修正为全局拟合+并行转换:

from dask_ml.preprocessing import StandardScaler

def elementwise_preprocess(tensor):
    # 逐张量执行预处理逻辑
    tensor[tensor < 1e-3] = 0.
    tensor[-1, ...] = 25. * tensor[-1, ...]
    tensor = tensor / 100.
    return tensor

# 对Dask Array做元素级并行预处理
X_processed = X_array.map_blocks(elementwise_preprocess, dtype=np.float32)

# 全局拟合StandardScaler(需展平特征维度)
scaler = StandardScaler()
# 展平为(n_samples, 2*125*125)
X_flat = X_processed.reshape((X_processed.shape[0], -1))
scaler.fit(X_flat)
# 转换后恢复原形状
X_scaled_flat = scaler.transform(X_flat)
X_scaled = X_scaled_flat.reshape(X_processed.shape)

3. TensorFlow分块训练整合

避免手动切片循环,利用Dask的延迟机制批量获取数据,适配TensorFlow输入要求:

import tensorflow as tf
from tensorflow.keras.applications import ResNet50
from tensorflow.keras.layers import Dense, GlobalAveragePooling2D
from tensorflow.keras.models import Model

# 定义适配输入的ResNet模型
base_model = ResNet50(
    weights=None,
    include_top=False,
    input_shape=(125, 125, 2)  # 适配转换后的通道维度
)
x = base_model.output
x = GlobalAveragePooling2D()(x)
predictions = Dense(1, activation='sigmoid')(x)  # 按需调整输出层
model = Model(inputs=base_model.input, outputs=predictions)
model.compile(optimizer='adam', loss='binary_crossentropy', metrics=['accuracy'])

# 批量遍历数据(以每个分区为一个批次,可调整batch_size)
for x_delayed, y_delayed in zip(X_scaled.to_delayed(), y_array.to_delayed()):
    # 计算当前批次的numpy数组
    x_batch = x_delayed.compute()
    # 转换形状为TensorFlow要求的(样本数, 高, 宽, 通道数)
    x_batch = np.transpose(x_batch, (0, 2, 3, 1))
    y_batch = y_delayed.compute()
    
    # 单批次训练
    model.train_on_batch(x_batch, y_batch)

更高效的批量迭代方式

使用iterbatches自定义批量大小:

# 生成批量迭代器
batch_iterator = da.map_blocks(
    lambda x, y: (x, y),
    X_scaled,
    y_array,
    meta=(np.float32, np.float32)
).iterbatches(batch_size=10)

for x_batch, y_batch in batch_iterator:
    x_batch = np.transpose(x_batch, (0, 2, 3, 1))
    model.train_on_batch(x_batch, y_batch)

核心注意事项

  • 分区合理性:保持split_row_groups=10的分区设置,避免单分区数据量过大触发OOM
  • 预处理一致性:必须全局拟合StandardScaler,禁止在单个分区单独拟合,否则会破坏数据分布
  • 输入形状适配:TensorFlow卷积层默认输入格式为(height, width, channels),需转换张量维度顺序
  • 最小化compute:仅在喂入模型时计算当前批次,禁止全局compute()加载全量数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 05:13:27