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

使用pandas+fastparquet写入Parquet时DATE逻辑类型配置问题

用pandas+fastparquet生成DATE逻辑类型(INT32物理类型)的Parquet文件并追加写入

问题背景

需要从数据库流式读取数据并追加写入同一Parquet文件,要求其中的日期列是DATE逻辑类型、INT32物理类型(该类型对应BigQuery的DATE列,类型不匹配会导致导入失败)。但直接写入字符串或转成datetime64[ns]类型,生成的Parquet列类型分别为STRING或Timestamp(INT64),均不符合要求。

解决方案

核心是利用fastparquet的自定义schema功能,将日期列指定为DATE类型,同时确保数据是datetime.date对象而非datetime64[ns]。

修改后的完整代码

import os
import pandas as pd
from sqlalchemy import create_engine
from sqlalchemy import text
from fastparquet import write

sql = "SELECT TO_VARCHAR(CURRENT_DATE, 'YYYYMMDD') AS \"REPORT_DATE\" FROM DUMMY;"

def create_stream_enabled_connection(username, password, host, port):
    conn_str = f"mydb://{username}:{password}@{host}:{port}"
    engine = create_engine(conn_str, connect_args={'encrypt': 'True', 'sslValidateCertificate':'False'})
    connection = engine.connect().execution_options(stream_results=True)
    return connection

connection = create_stream_enabled_connection(username, password, host, port)

ROWS_IN_CHUNK = 500000
local_path = "target_file.parquet"

# 自定义schema:指定REPORT_DATE为Parquet DATE类型
custom_schema = {"REPORT_DATE": "date"}

# 流式读取并写入Parquet
for idx, dataframe_chunk in enumerate(pd.read_sql(text(sql), connection, chunksize=ROWS_IN_CHUNK)):
    # 将字符串日期转换为datetime.date对象(fastparquet识别DATE类型的前提)
    dataframe_chunk["REPORT_DATE"] = pd.to_datetime(dataframe_chunk["REPORT_DATE"], format="%Y%m%d").dt.date
    
    if idx == 0 and not os.path.exists(local_path):
        # 首次写入:指定schema
        write(
            local_path,
            dataframe_chunk,
            index=False,
            schema=custom_schema,
            engine="fastparquet"
        )
    else:
        # 追加写入:必须复用相同schema保证类型一致性
        write(
            local_path,
            dataframe_chunk,
            index=False,
            append=True,
            schema=custom_schema,
            engine="fastparquet"
        )

关键说明

  1. 数据类型转换:pd.to_datetime(...).dt.date将字符串日期转为datetime.date对象,这是fastparquet识别为DATE类型的必要条件。
  2. 自定义schema:通过schema={"REPORT_DATE": "date"}明确指定列类型为Parquet的DATE逻辑类型,其物理存储自动为INT32(存储从1970-01-01开始的天数偏移)。
  3. 追加一致性:追加写入时必须使用相同的schema,避免因类型不匹配导致写入失败。

验证方法

用fastparquet读取生成的文件,检查列类型:

from fastparquet import ParquetFile
pf = ParquetFile(local_path)
print(pf.schema)

输出中REPORT_DATE的类型应为date,对应物理存储为INT32。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 20:41:27