使用AWS Wrangler COPY时Redshift与S3 Parquet数据类型不匹配问题
Redshift COPY 因Parquet Schema不兼容报错的解决办法
问题背景
我们每小时运行大量Lambda任务,从API拉取数据转换为DataFrame后,用method='multi'批量插入单节点Redshift集群。但频繁的插入导致集群CPU占用率100%,其他任务超时。于是切换测试AWS Wrangler的Redshift COPY方案:先将DataFrame导出为Parquet文件存储到S3,再通过wr.redshift.copy()追加到Redshift表。
为验证UNLOAD和COPY流程,我们先通过wr.redshift.unload()从样本表导出少量数据生成DataFrame,再尝试导入同结构的测试表,却触发错误:
Spectrum Scan Error: <s3path> has an incompatible Parquet schema for column
具体表现为TIMESTAMP列被转为CHAR类型、INT2列被转为INT64类型,与目标表结构不匹配。
相关代码
UNLOAD 导出代码
test_df = wr.redshift.unload( sql=f"SELECT * FROM {schema}.{table_name};", path="s3://bucket-name/copy_test/", keep_files=False, con=rs_con )
COPY 导入代码
wr.redshift.copy( df=test_df, table=table_name, schema=schema, keep_files=False, path='s3://bucket-name/copy_test/', use_column_names=True, index=False, con=rs_con )
问题根源
Redshift UNLOAD导出Parquet时,会将原生类型转换为Parquet兼容类型,但Pandas读取Parquet时又会转换为自身默认类型,导致类型错位:
- Redshift的
TIMESTAMP被UNLOAD默认以文本格式存储为Parquet的STRING,Pandas读取后变为object/CHAR类型 - Redshift的
INT2(SMALLINT)因Parquet无对应类型,被UNLOAD转为INT64,Pandas读取后保持该类型,与目标表的INT2不兼容
解决步骤
1. 优化UNLOAD的类型映射
修改UNLOAD调用,通过unload_kwargs强制指定时间类型为Parquet的TIMESTAMP,确保数值类型正确转换:
test_df = wr.redshift.unload( sql=f"SELECT * FROM {schema}.{table_name};", path="s3://bucket-name/copy_test/", keep_files=False, con=rs_con, unload_kwargs={ "FORMAT": "PARQUET", "PARQUET_COMPRESSION": "SNAPPY", "TIMEFORMAT": "epochmillisecs" # 将时间转为毫秒时间戳,UNLOAD会存储为Parquet TIMESTAMP类型 } )
2. 手动对齐DataFrame列类型
如果UNLOAD后Pandas的列类型仍不符合要求,手动转换对应列:
import pandas as pd # 转换TIMESTAMP列 test_df['your_timestamp_column'] = pd.to_datetime(test_df['your_timestamp_column']) # 转换INT2列(注意:需确保数值在-32768到32767范围内,避免溢出) test_df['your_int2_column'] = test_df['your_int2_column'].astype('int16')
3. COPY时显式指定类型映射
在COPY调用中通过copy_kwargs明确告诉Redshift如何解析Parquet列的类型,与目标表结构对齐:
wr.redshift.copy( df=test_df, table=table_name, schema=schema, keep_files=False, path='s3://bucket-name/copy_test/', use_column_names=True, index=False, con=rs_con, copy_kwargs={ "FORMAT": "PARQUET", # 按目标表结构列出列名和对应类型 "COLUMNS": "(your_timestamp_column TIMESTAMP, your_int2_column SMALLINT, other_col VARCHAR(255))" } )
额外优化建议
- COPY方案本身比
method='multi'的批量插入高效得多,调整后单节点Redshift的CPU占用率会明显下降 - 导出Parquet时使用SNAPPY压缩,减少S3存储占用和数据传输时间
- 控制每个Lambda任务处理的DataFrame大小,避免单批次数据过大给集群带来压力
内容的提问来源于stack exchange,提问作者lordbuckshot
相关产品推荐
相关产品推荐

