Redshift Glue Job Bigint字段插入失败问题求助(Python3+Spark2.2)
解决Glue Job向Redshift插入bigint字段失败的问题
我来帮你排查这个bigint字段插入失败的问题,结合你的代码和场景,可能有几个常见的原因,咱们一步步来解决:
1. 先确认数据源的selling字段实际类型
你在ApplyMapping里指定了("selling", "bigint", "selling", "bigint"),但首先要核实Glue Catalog中test.tbl_test的selling字段是不是真的是bigint类型。有时候Glue爬虫会误识别数据源(比如S3的CSV/Parquet文件)的字段类型,比如把数值识别成string或者decimal,这样后续的映射就会出错。
- 检查方式:登录Glue控制台,找到对应的数据库和表,查看表结构里的
selling字段类型。 - 修复方案:如果类型不符,要么重新运行Glue爬虫正确识别字段类型,要么在
ApplyMapping里手动做类型转换,比如原类型是string的话:applymapping1 = ApplyMapping.apply(frame = datasource0, mappings = [ ("testdata", "string", "testdata", "string"), ("selling", "string", "selling", "bigint") ], transformation_ctx = "applymapping1")
2. 优化DynamicFrame的类型解析逻辑
你用了ResolveChoice.apply(choice = "make_cols"),这个策略会生成多个同名不同类型的字段(比如selling_int、selling_string),反而会导致最终写入Redshift时字段不匹配。
建议改成明确指定类型的策略,比如强制把selling字段转成bigint:
resolvechoice2 = ResolveChoice.apply( frame = applymapping1, specs = [("selling", "cast:bigint")], # 针对单个字段做类型强转 transformation_ctx = "resolvechoice2" )
3. 核对Redshift侧的配置
- 确认Redshift目标表
tbl_test的selling字段确实是bigint类型(Redshift的bigint是8字节整数,范围是-9223372036854775808到9223372036854775807,如果你的数据超出这个范围也会插入失败)。 - 检查Glue连接使用的IAM角色权限:需要具备Redshift目标表的读写权限,以及
args["TempDir"]对应S3路径的读写权限(Glue写Redshift是先把数据写到S3临时目录,再通过COPY命令导入)。
4. 通过日志定位具体错误
很多时候直接看Job的运行日志就能找到问题根源:
- 进入Glue控制台,找到对应的Job,查看运行日志里的报错信息,比如是类型转换失败、数据超出范围,还是权限不足。
调整后的完整代码示例
整合了上述优化点的代码:
import sys from awsglue.transforms import * from awsglue.utils import getResolvedOptions from pyspark.context import SparkContext from awsglue.context import GlueContext from awsglue.job import Job args = getResolvedOptions(sys.argv, ['TempDir','JOB_NAME']) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args['JOB_NAME'], args) # 读取数据源 datasource0 = glueContext.create_dynamic_frame.from_catalog( database = "test", table_name = "tbl_test", transformation_ctx = "datasource0" ) # 字段映射:根据实际数据源类型调整转换逻辑 applymapping1 = ApplyMapping.apply( frame = datasource0, mappings = [ ("testdata", "string", "testdata", "string"), ("selling", "bigint", "selling", "bigint") ], transformation_ctx = "applymapping1" ) # 明确解析selling字段为bigint,避免类型歧义 resolvechoice2 = ResolveChoice.apply( frame = applymapping1, specs = [("selling", "cast:bigint")], transformation_ctx = "resolvechoice2" ) # 过滤空值 dropnullfields3 = DropNullFields.apply( frame = resolvechoice2, transformation_ctx = "dropnullfields3" ) # 写入Redshift datasink4 = glueContext.write_dynamic_frame.from_jdbc_conf( frame = dropnullfields3, catalog_connection = "redshift_con", connection_options = {"dbtable": "tbl_test", "database": "test"}, redshift_tmp_dir = args["TempDir"], transformation_ctx = "datasink4" ) job.commit()
内容的提问来源于stack exchange,提问作者user2768132
相关产品推荐
相关产品推荐

