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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 22:05:17