将Parquet文件导入Redshift遇类型不兼容错误,如何解决?
Redshift导入Parquet文件Schema兼容问题解决
问题背景
通过SQLAlchemy在Redshift创建表:
from sqlalchemy.ext.declarative import declarative_base from sqlalchemy import Column from sqlalchemy_redshift.dialect import INTEGER, VARCHAR Base = declarative_base() class SomeTable(Base): __tablename__ = "some_table" id = Column(INTEGER, primary_key=True) some_column = Column(VARCHAR(128)) # 连接Redshift并创建表 Base.metadata.create_all(engine)
生成的Redshift表结构:
CREATE TABLE public.some_table ( id integer NOT NULL ENCODE az64, some_column character varying(128) ENCODE lzo, ) DISTSTYLE AUTO SORTKEY ( id );
将pandas DataFrame(id为int64类型)导出为Parquet并上传S3,使用COPY命令导入时出现错误:
psycopg2.errors.InternalError_: Spectrum Scan Error DETAIL: ----------------------------------------------- error: Spectrum Scan Error code: 15007 context: File 'https://s3.eu-central-1.amazonaws.com/bucket/1681141668/some_table.parquet' has an incompatible Parquet schema for column 's3://bucket/1681141668/some_table.parquet.id'. Column type: INT, Parquet schema: optional int query: 2627 location: dory_util.cpp:1450 process: worker_thread [pid=19137] -----------------------------------------------
问题原因
Redshift表中id列是非空32位整数(integer NOT NULL),但生成的Parquet文件中id是可选64位整数(INT64,允许NULL),两者在整数位数和可空性上不匹配,导致Spectrum扫描时类型兼容错误。
解决方案
方法一:调整Parquet文件Schema(推荐)
修改Parquet生成代码,将id列设置为非空32位整数,完全匹配Redshift表的定义:
from tempfile import NamedTemporaryFile from pyarrow import Table, int32, schema, string, field from pyarrow.parquet import write_table with NamedTemporaryFile() as file: parquet_table = Table.from_pandas( df, schema=schema( [ # 定义为非空32位整数,对应Redshift的integer类型 field("id", int32(), nullable=False), field("some_column", string(), nullable=True), ] ), ) write_table(parquet_table, file) # 上传文件至S3的代码...
修改后用parquet-tools inspect验证:
id列的physical_type应为INT32max_definition_level应为0(表示列不可为空)
方法二:修改Redshift表结构(不推荐)
若无法修改Parquet文件,可将Redshift表的id列类型改为bigint NOT NULL(匹配Parquet的INT64),但需重建表迁移数据:
-- 创建新表 CREATE TABLE public.some_table_new ( id bigint NOT NULL ENCODE az64, some_column character varying(128) ENCODE lzo, ) DISTSTYLE AUTO SORTKEY ( id ); -- 迁移现有数据 INSERT INTO public.some_table_new SELECT * FROM public.some_table; -- 替换原表 DROP TABLE public.some_table; ALTER TABLE public.some_table_new RENAME TO some_table;
此方法需注意数据迁移成本,且主键类型变更可能影响关联业务,仅作为备选方案。
内容的提问来源于stack exchange,提问作者doublethink13
相关产品推荐
相关产品推荐

