使用pyarrow.dataset Scanner将CSV转Parquet时如何修改列类型
解决PyArrow Dataset解析CSV时区时间并转Parquet的问题
核心问题原因
直接给ds.dataset传入schema失败,是因为PyArrow的CSV解析器默认无法识别带时区后缀的时间格式,必须显式指定时间解析规则,同时匹配正确的带时区timestamp类型。
分步解决方案
1. 定义包含正确时间类型的Schema
原时间格式2022-01-03 14:30:00.000652 UTC是微秒精度+UTC时区的时间,所以Schema中time列要定义为带时区的timestamp类型:
import pyarrow as pa import pyarrow.dataset as ds source = "foo.csv" dest = "Data/parquet" # 按实际CSV列定义完整Schema,这里仅突出time列 schema = pa.schema([ ("time", pa.timestamp("us", tz="UTC")), # 其他列示例:("col1", pa.int64()), ("col2", pa.string()), ... ])
2. 配置CSV解析选项,指定时间格式
通过CsvParseOptions传入时间解析模板,完全匹配原字符串的格式:
csv_parse_opts = ds.CsvParseOptions( timestamp_parsers=["%Y-%m-%d %H:%M:%S.%f UTC"] )
%f匹配微秒部分(000652),UTC匹配后缀的时区标识,必须和原字符串格式完全一致。
3. 创建可正确解析时间的CSV Dataset
将Schema和解析选项传入ds.dataset,确保CSV被正确解析为目标类型:
csv_ds = ds.dataset( source, format="csv", schema=schema, parse_options=csv_parse_opts )
4. (可选)添加分区字段(从time列提取date)
如果需要按date分区但原CSV无该列,可在Scanner中通过计算生成date列:
# 扫描所有原始列,同时添加从time转换来的date列 scanner = csv_ds.scanner( columns=["*", ds.field("time").cast(pa.date32()).alias("date")] )
5. 写入Parquet文件
最后用配置好的Scanner写入,指定分区规则:
partitioning = ds.partitioning(pa.schema([("date", pa.date32())]), flavor="hive") ds.write_dataset( scanner, dest, format="parquet", partitioning=partitioning, existing_data_behavior="overwrite" # 根据需求选择保留或覆盖已有数据 )
关键注意事项
- 若CSV中有其他列,Schema必须完整定义所有列的类型,否则PyArrow会尝试自动推断,可能出现类型不匹配。
- 时区必须明确指定(这里是UTC),否则解析后的timestamp会是无时区类型,后续处理易出问题。
timestamp_parsers的格式字符串必须和原始时间字符串完全匹配,包括空格、时区后缀等细节。
内容的提问来源于stack exchange,提问作者hjs
相关产品推荐
相关产品推荐

