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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 17:23:25