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错误,错误信息为:
错误发生在读取批次阶段,还未执行到后续的cast逻辑。pyarrow.lib.ArrowInvalid: In CSV column #49: CSV conversion error to null: invalid value '0.0000' - 若移除读取批次后转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
相关产品推荐
相关产品推荐

