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

Pyarrow/Parquet批量处理时将null列强制转换为string类型报错求解

问题背景

在内存受限系统中,需要将解压后约700MB的tar.gz压缩CSV文件转换为Parquet格式,因此采用分批次处理方案。目前已实现流式读取tar.gz压缩包、提取目标CSV文件、通过pyarrow的open_csv()接口分块读取数据的逻辑,计划分批将数据写入Parquet文件。

故障现象

目标CSV存在大量前序行全为空、仅在约第50万行左右才出现有效值的列,pyarrow初始识别这类列的dtype为null类型。手动修改Schema将这类null类型列设置为string以兼容后续有效值时,出现两类报错:

  • 读取批次后将批次转为Table再cast到修改后目标Schema的流程中,调用batch = reader.read_next_batch()环节就抛出ArrowInvalid错误,错误信息为:
    pyarrow.lib.ArrowInvalid: In CSV column #49: CSV conversion error to null: invalid value '0.0000'
    
    错误发生在读取批次阶段,还未执行到后续的cast逻辑。
  • 若移除读取批次后转Table执行cast的逻辑,直接写入读取到的批次时,会抛出Schema不匹配错误:读取到的批次中对应列仍为null类型,与ParquetWriter初始化时指定的string类型Schema不符。
原有实现代码
import io
import os
import tarfile
import pyarrow as pa
import pyarrow.parquet as pq
import pyarrow.csv as csv
import logging

srcs = list()
path = "C:\\data"
for root, dirs, files in os.walk(path):
   for name in files:
       if name.endswith("tar.gz"):
           srcs.append(os.path.join(root, name))

for source_file_name in srcs:

    file_name: str = source_file_name.replace(".tar.gz", "")
    target_file_name: str = source_file_name.replace(".tar.gz", ".parquet")
    clean_file_name: str = os.path.basename(source_file_name.replace(".tar.gz", ""))

    # 处理CSV文件 保留目录结构
    logging.info(f"Processing '{source_file_name}'.")
    with io.open(source_file_name, "rb") as file_obj_in:
        file_obj_in.seek(0)
        with tarfile.open(fileobj=file_obj_in, mode="r") as tf:
            file_obj = tf.extractfile(f"{clean_file_name}.csv")
            file_obj.seek(0)
            reader = csv.open_csv(file_obj, read_options=csv.ReadOptions(block_size=25*1024*1024))
            schema = reader.schema
            null_cols = list()
            for index, entry in enumerate(schema.types):
                if entry.equals(pa.null()):
                    schema = schema.set(index, schema.field(index).with_type(pa.string()))
                    null_cols.append(index)
            
            with pq.ParquetWriter(target_file_name, schema) as writer:
                while True:
                    try:
                        batch = reader.read_next_batch()
                        table = pa.Table.from_batches(batches=[batch]).cast(target_schema=schema)
                        batch = table.to_batches()[0]
                        writer.write_batch(batch)
                    except StopIteration:
                        break
故障原因

csv.open_csv()初始化完成时,列解析规则就已经固定。初始化后手动修改返回的schema对象,不会改变reader内部的解析逻辑:被自动推断为null类型的列,reader会始终按null类型做解析校验,读到非空值时直接在读取阶段抛出错误,后续的cast逻辑根本没有执行机会。

解决方案

核心思路是在CSV reader初始化前就显式指定所有列的目标类型,让reader从解析第一个块开始就按指定类型处理数据,从根源避免null类型解析报错,同时省掉后续cast的内存开销。
具体操作步骤:

  • 第一次打开CSV文件时,仅做初始schema推断,识别出所有被自动判定为null类型的列,将这些列的类型统一覆盖为string,生成最终目标schema。
  • 将CSV文件指针重置到起始位置,重新初始化CSV reader,通过csv.ConvertOptions的column_types参数传入提前定义好的列类型映射,让reader按指定类型解析所有列。
  • 直接将读取到的批次写入Parquet文件即可,不需要额外做类型转换。

修正后的可运行代码:

import io
import os
import tarfile
import pyarrow as pa
import pyarrow.parquet as pq
import pyarrow.csv as csv
import logging

srcs = list()
path = "C:\\data"
for root, dirs, files in os.walk(path):
    for name in files:
        if name.endswith("tar.gz"):
            srcs.append(os.path.join(root, name))

# 固定分批读取块大小 25MB
BLOCK_SIZE = 25 * 1024 * 1024

for source_file_name in srcs:
    target_file_name: str = source_file_name.replace(".tar.gz", ".parquet")
    clean_file_name: str = os.path.basename(source_file_name.replace(".tar.gz", ""))

    logging.info(f"Processing '{source_file_name}'.")
    with io.open(source_file_name, "rb") as file_obj_in:
        file_obj_in.seek(0)
        with tarfile.open(fileobj=file_obj_in, mode="r") as tf:
            file_obj = tf.extractfile(f"{clean_file_name}.csv")
            # 第一次打开reader 仅用于推断schema
            pre_reader = csv.open_csv(
                file_obj,
                read_options=csv.ReadOptions(block_size=BLOCK_SIZE)
            )
            target_schema = pre_reader.schema
            # 构建列类型覆盖映射
            column_type_overrides = {}
            for idx, field in enumerate(target_schema):
                if field.type.equals(pa.null()):
                    # null列统一覆盖为string类型
                    new_field = field.with_type(pa.string())
                    target_schema = target_schema.set(idx, new_field)
                    column_type_overrides[field.name] = pa.string()
            # 关闭预读reader 重置文件指针到CSV开头
            pre_reader.close()
            file_obj.seek(0)

            # 重新初始化reader 传入显式类型配置
            convert_options = csv.ConvertOptions(column_types=column_type_overrides)
            reader = csv.open_csv(
                file_obj,
                read_options=csv.ReadOptions(block_size=BLOCK_SIZE),
                convert_options=convert_options
            )

            # 分批写入Parquet
            with pq.ParquetWriter(target_file_name, target_schema) as writer:
                while True:
                    try:
                        batch = reader.read_next_batch()
                        writer.write_batch(batch)
                    except StopIteration:
                        break
            reader.close()

说明:该方案全程仅占用单批次25MB左右的内存,完全适配内存受限场景,不需要额外落地中间文件,也不需要对CSV做全量预扫描。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 16:48:28