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

如何将SQLAlchemy/SQL Server数据类型转换为PyArrow数据类型?

我正在开发一个从SQL Server每日加载数据到Parquet文件的程序,同时也支持导出为CSV或JSON格式。我使用pyodbc、SQLAlchemy、pandas和pyarrow来完成这项工作。

import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq
import sys
import argparse
import sqlalchemy as sa
import os
...
connection_string = f'mssql+pyodbc://{server}/{database}?driver=ODBC+Driver+17+for+SQL+Server&trusted_connection=yes&read_only=true'
engine = sa.create_engine(connection_string)
...
sql = sa.text(f"""
    SELECT MAIN.*
        , etl_dest.ETL_PROCESS_EXECUTION_IDENTIFIER as [ETL_EVENT_PROCESS_EXECUTION_IDENTIFIER]
    FROM {database}.{schema}.{table} MAIN
    INNER JOIN {database}.audit.ETL_EVENT_DESTINATION etl_dest on MAIN.ETL_EVENT_DESTINATION_IDENTIFIER = etl_dest.ETL_EVENT_DESTINATION_IDENTIFIER  
    INNER JOIN {database}.Audit.ETL_PROCESS_EXECUTION etl_proc on etl_dest.ETL_PROCESS_EXECUTION_IDENTIFIER = etl_proc.ETL_PROCESS_EXECUTION_IDENTIFIER  
    INNER JOIN {database}.Audit.ETL_BATCH_EXECUTION etl_batch on etl_proc.ETL_BATCH_EXECUTION_IDENTIFIER = etl_batch.ETL_BATCH_EXECUTION_IDENTIFIER  
    WHERE
        etl_proc.ETL_BATCH_EXECUTION_IDENTIFIER = :batch_identifier
    """)
使用SQLAlchemy检索表架构
# inspector = sa.inspect(engine)
# columns_info = inspector.get_columns(table, schema=schema)
# columns_info.append({'name': 'RWB_ETL_EVENT_PROCESS_EXECUTION_IDENTIFIER', 'type': sa.INTEGER(), 'nullable': False, 'default': None, 'autoincrement': False, 'comment': None})

# 建立与主数据库的连接
with engine.connect().execution_options(stream_results=True) as connection:
    # 执行查询并分块获取结果
    # TODO: 如果有人实现SQLAlchemy到PyArrow类型的转换,这里指定dtype=可能会提升性能
    chunks = pd.read_sql(sql, connection, params={"batch_identifier": batch_identifier}, chunksize=10000, dtype_backend='pyarrow')
    # 将分块数据写入输出文件
    if file_format == 'parquet':
        writer = None
        for chunk in chunks:
            batch = pa.RecordBatch.from_pandas(chunk)
            if writer is None:
                writer = pq.ParquetWriter(output_file, batch.schema)
            writer.write_batch(batch)
        writer.close()
    elif file_format == 'csv':
        header = True
        for chunk in chunks:
            chunk.to_csv(output_file, index=False, sep=delimiter, mode='a', header=header)
            header = False
    elif file_format == 'json':
        for chunk in chunks:
            chunk.to_json(output_file, orient='records', lines=True, force_ascii=False, mode='a')
    else:
        raise ValueError(f"不支持的文件格式: {file_format}")
  • 核心需求:获取SQL Server表对应的PyArrow数据类型定义,numpy数据类型也可接受
  • 用途:优化pandas从SQL Server读取数据的性能,保证各分块数据类型一致
  • 问题背景:若不设置dtype_backend='pyarrow',写入Parquet文件时会出现文件与表架构不匹配的问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 01:57:17