使用AWS Glue迁移S3 CSV数据到Redshift时pyWriteDynamicFrame报错
尝试使用AWS Glue将S3存储桶中的CSV数据迁移至Amazon Redshift表时,任务执行失败,报错信息如下:
Error Category: UNCLASSIFIED_ERROR; An error occurred while calling o209.pyWriteDynamicFrame. Exception thrown in awaitResult:
任务失败后,Redshift表的Schema已创建,但未导入任何数据。已检查Glue日志和Redshift的stl_load_errors表,均未发现有效错误线索。
Glue任务脚本如下:
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 from awsglue import DynamicFrame args = getResolvedOptions(sys.argv, ["JOB_NAME"]) sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session job = Job(glueContext) job.init(args["JOB_NAME"], args) # Script generated for node Amazon S3 AmazonS3_node1704977097218 = glueContext.create_dynamic_frame.from_options( format_options={"quoteChar": '"', "withHeader": True, "separator": ","}, connection_type="s3", format="csv", connection_options={"paths": ["s3://smedia-data-raw-dev/google/"], "recurse": True}, transformation_ctx="AmazonS3_node1704977097218", ) # Script generated for node Change Schema ChangeSchema_node1705498745919 = ApplyMapping.apply( frame=AmazonS3_node1704977097218, mappings=[ ("resourcename", "string", "resourcename", "varchar"), ("status", "string", "status", "varchar"), ("basecampaign", "string", "basecampaign", "varchar"), ("name", "string", "name", "varchar"), ("id", "string", "id", "varchar"), ("campaignbudget", "string", "campaignbudget", "varchar"), ("startdate", "string", "startdate", "varchar"), ("enddate", "string", "enddate", "varchar"), ( "adservingoptimizationstatus", "string", "adservingoptimizationstatus", "varchar", ), ("advertisingchanneltype", "string", "advertisingchanneltype", "varchar"), ("advertisingchannelsubtype", "string", "advertisingchannelsubtype", "varchar"), ("experimenttype", "string", "experimenttype", "varchar"), ("servingstatus", "string", "servingstatus", "varchar"), ("biddingstrategytype", "string", "biddingstrategytype", "varchar"), ("domainname", "string", "domainname", "varchar"), ("languagecode", "string", "languagecode", "varchar"), ("usesuppliedurlsonly", "string", "usesuppliedurlsonly", "varchar"), ("positivegeotargettype", "string", "positivegeotargettype", "varchar"), ("negativegeotargettype", "string", "negativegeotargettype", "varchar"), ("paymentmode", "string", "paymentmode", "varchar"), ("optimizationgoaltypes", "string", "optimizationgoaltypes", "varchar"), ("date", "string", "date", "varchar"), ("averagecost", "string", "averagecost", "varchar"), ("clicks", "string", "clicks", "varchar"), ("costmicros", "string", "costmicros", "varchar"), ("impressions", "string", "impressions", "varchar"), ("useaudiencegrouped", "string", "useaudiencegrouped", "varchar"), ( "activeviewmeasurablecostmicros", "string", "activeviewmeasurablecostmicros", "varchar", ), ("costperallconversions", "string", "costperallconversions", "varchar"), ("costperconversion", "string", "costperconversion", "varchar"), ("invalidclicks", "string", "invalidclicks", "varchar"), ("publisherpurchasedclicks", "string", "publisherpurchasedclicks", "varchar"), ("averagepageviews", "string", "averagepageviews", "varchar"), ("videoviews", "string", "videoviews", "varchar"), ( "allconversionsbyconversiondate", "string", "allconversionsbyconversiondate", "varchar", ), ( "allconversionsvaluebyconversiondate", "string", "allconversionsvaluebyconversiondate", "varchar", ), ( "conversionsbyconversiondate", "string", "conversionsbyconversiondate", "varchar", ), ( "conversionsvaluebyconversiondate", "string", "conversionsvaluebyconversiondate", "varchar", ), ( "valueperallconversionsbyconversiondate", "string", "valueperallconversionsbyconversiondate", "varchar", ), ( "valueperconversionsbyconversiondate", "string", "valueperconversionsbyconversiondate", "varchar", ), ("allconversions", "string", "allconversions", "varchar"), ( "absolutetopimpressionpercentage", "string", "absolutetopimpressionpercentage", "varchar", ), ( "searchabsolutetopimpressionshare", "string", "searchabsolutetopimpressionshare", "varchar", ), ("averagecpc", "string", "averagecpc", "varchar"), ("searchimpressionshare", "string", "searchimpressionshare", "varchar"), ("searchtopimpressionshare", "string", "searchtopimpressionshare", "varchar"), ("activeviewctr", "string", "activeviewctr", "varchar"), ("ctr", "string", "ctr", "varchar"), ("relativectr", "string", "relativectr", "varchar"), ], transformation_ctx="ChangeSchema_node1705498745919", ) # Script generated for node Amazon Redshift AmazonRedshift_node1705498761115 = glueContext.write_dynamic_frame.from_options( frame=ChangeSchema_node1705498745919, connection_type="redshift", connection_options={ "redshiftTmpDir": "s3://aws-glue-assets-191965435652-us-east-1/temporary/", "useConnectionProperties": "true", "dbtable": "public.google_test_2", "connectionName": "Redshift connection", "preactions": "DROP TABLE IF EXISTS public.google_test_2; CREATE TABLE IF NOT EXISTS public.google_test_2 (resourcename VARCHAR, status VARCHAR, basecampaign VARCHAR, name VARCHAR, id VARCHAR, campaignbudget VARCHAR, startdate VARCHAR, enddate VARCHAR, adservingoptimizationstatus VARCHAR, advertisingchanneltype VARCHAR, advertisingchannelsubtype VARCHAR, experimenttype VARCHAR, servingstatus VARCHAR, biddingstrategytype VARCHAR, domainname VARCHAR, languagecode VARCHAR, usesuppliedurlsonly VARCHAR, positivegeotargettype VARCHAR, negativegeotargettype VARCHAR, paymentmode VARCHAR, optimizationgoaltypes VARCHAR, date VARCHAR, averagecost VARCHAR, clicks VARCHAR, costmicros VARCHAR, impressions VARCHAR, useaudiencegrouped VARCHAR, activeviewmeasurablecostmicros VARCHAR, costperallconversions VARCHAR, costperconversion VARCHAR, invalidclicks VARCHAR, publisherpurchasedclicks VARCHAR, averagepageviews VARCHAR, videoviews VARCHAR, allconversionsbyconversiondate VARCHAR, allconversionsvaluebyconversiondate VARCHAR, conversionsbyconversiondate VARCHAR, conversionsvaluebyconversiondate VARCHAR, valueperallconversionsbyconversiondate VARCHAR, valueperconversionsbyconversiondate VARCHAR, allconversions VARCHAR, absolutetopimpressionpercentage VARCHAR, searchabsolutetopimpressionshare VARCHAR, averagecpc VARCHAR, searchimpressionshare VARCHAR, searchtopimpressionshare VARCHAR, activeviewctr VARCHAR, ctr VARCHAR, relativectr VARCHAR);", }, transformation_ctx="AmazonRedshift_node1705498761115", ) job.commit()
1. 修复Redshift表字段长度问题
Redshift中VARCHAR类型未指定长度时,默认仅支持1个字符。你的preactions中所有字段均使用无长度的VARCHAR,当数据内容超过1字符时会导致插入失败,但此类错误可能未被stl_load_errors捕获。
解决方案:
修改preactions中的CREATE语句,为每个VARCHAR字段指定合理长度,示例如下:
CREATE TABLE IF NOT EXISTS public.google_test_2 ( resourcename VARCHAR(255), status VARCHAR(50), basecampaign VARCHAR(255), name VARCHAR(255), id VARCHAR(100), campaignbudget VARCHAR(100), startdate VARCHAR(20), enddate VARCHAR(20), adservingoptimizationstatus VARCHAR(100), advertisingchanneltype VARCHAR(100), advertisingchannelsubtype VARCHAR(100), experimenttype VARCHAR(100), servingstatus VARCHAR(100), biddingstrategytype VARCHAR(100), domainname VARCHAR(255), languagecode VARCHAR(20), usesuppliedurlsonly VARCHAR(20), positivegeotargettype VARCHAR(100), negativegeotargettype VARCHAR(100), paymentmode VARCHAR(100), optimizationgoaltypes VARCHAR(255), date VARCHAR(20), averagecost VARCHAR(50), clicks VARCHAR(50), costmicros VARCHAR(50), impressions VARCHAR(50), useaudiencegrouped VARCHAR(20), activeviewmeasurablecostmicros VARCHAR(50), costperallconversions VARCHAR(50), costperconversion VARCHAR(50), invalidclicks VARCHAR(50), publisherpurchasedclicks VARCHAR(50), averagepageviews VARCHAR(50), videoviews VARCHAR(50), allconversionsbyconversiondate VARCHAR(50), allconversionsvaluebyconversiondate VARCHAR(50), conversionsbyconversiondate VARCHAR(50), conversionsvaluebyconversiondate VARCHAR(50), valueperallconversionsbyconversiondate VARCHAR(50), valueperconversionsbyconversiondate VARCHAR(50), allconversions VARCHAR(50), absolutetopimpressionpercentage VARCHAR(50), searchabsolutetopimpressionshare VARCHAR(50), averagecpc VARCHAR(50), searchimpressionshare VARCHAR(50), searchtopimpressionshare VARCHAR(50), activeviewctr VARCHAR(50), ctr VARCHAR(50), relativectr VARCHAR(50) );
2. 验证临时S3目录权限
Glue写入Redshift时,会先将数据转存到redshiftTmpDir指定的S3路径,再通过Redshift的COPY命令加载:
- 确认Glue执行角色对
s3://aws-glue-assets-191965435652-us-east-1/temporary/拥有s3:GetObject、s3:PutObject、s3:ListBucket权限 - 确认Redshift集群关联的IAM角色(用于COPY操作)对该临时目录拥有相同的读写权限
- 检查Redshift集群的安全组是否允许Glue的IP范围访问(或配置VPC端点实现内网访问)
3. 简化preactions逻辑
当前preactions先执行DROP再CREATE,可能引发锁表或执行顺序问题。建议改为仅保留CREATE TABLE IF NOT EXISTS,避免每次执行都删除表:
"preactions": "CREATE TABLE IF NOT EXISTS public.google_test_2 (resourcename VARCHAR(255), status VARCHAR(50), ...);"
4. 确认数据读取有效性
在脚本中添加数据计数,验证Glue是否成功读取S3中的CSV数据:
# 在读取S3数据的节点后添加 print(f"读取到的记录数: {AmazonS3_node1704977097218.count()}")
如果计数为0,说明S3路径下无有效CSV文件(或路径错误、权限不足),需检查S3路径是否正确、文件是否存在、Glue角色是否有S3读取权限。
5. 启用详细日志排查
在Glue任务配置中,将日志级别调整为DEBUG,同时在脚本开头添加日志配置:
sc.setLogLevel("DEBUG")
执行任务后,查看CloudWatch日志中的详细报错信息,通常能找到awaitResult内部的具体异常。
6. 检查Redshift连接配置
- 确认Glue中配置的"Redshift connection"名称完全匹配,且连接关联的IAM角色拥有访问Redshift集群的权限
- 验证Redshift集群是否处于可用状态,端点地址正确,且Glue能通过网络访问集群
内容的提问来源于stack exchange,提问作者Matt J

