如何将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
相关产品推荐
相关产品推荐

