AWS Glue Job触发SchemaColumnConvertNotSupportedException问题求助
问题描述
在AWS Glue中读取S3存储的Parquet文件创建DynamicFrame后,调用show()方法时抛出SchemaColumnConvertNotSupportedException,提示Parquet列无法转换。尝试使用ResolveChoice将NumberOfEmployees列转为int类型,但问题依旧;实际需求是将该列从int转为String,但转换未生效。
相关代码
datasource0 = glueContext.create_dynamic_frame.from_options( format_options={}, connection_type="s3", format="parquet", connection_options={ "paths": s3_path_lst, "recurse": True, }, transformation_ctx="datasource0", ) datasource0.printSchema() datasource_dyf = datasource0.resolveChoice(specs=[('NumberOfEmployees', 'cast:int')]) datasource_dyf.printSchema() datasource_dyf.show()
异常信息
org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1996) at org.apache.spark.api.python.BasePythonRunner$WriterThread.run(PythonRunner.scala:232) Caused by: org.apache.spark.sql.execution.datasources.SchemaColumnConvertNotSupportedException at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.constructConvertNotSupportedException(VectorizedColumnReader.java:339) at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readIntBatch(VectorizedColumnReader.java:571) at org.apache.spark.sql.execution.datasources.parquet.VectorizedColumnReader.readBatch(VectorizedColumnReader.java:294) at org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextBatch(VectorizedParquetRecordReader.java:295) at org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader.nextKeyValue(VectorizedParquetRecordReader.java:196) at org.apache.spark.sql.execution.datasources.RecordReaderIterator.hasNext(RecordReaderIterator.scala:37) at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:159) at
解决方案
核心问题分析
- 代码逻辑与需求不符:实际要将
NumberOfEmployees从int转String,但代码中写的是cast:int,相当于尝试将非int类型转int,反而触发转换异常。 - Spark向量化读取提前校验:Glue读取Parquet时,Spark向量化读取器会在
ResolveChoice转换前校验列类型,若Parquet文件中该列存在无法匹配目标类型的数值(如超大整数、非整数格式),会直接抛出异常。
具体解决步骤
步骤1:修正ResolveChoice转换规则
将转换规则改为cast:string,匹配实际需求:
datasource_dyf = datasource0.resolveChoice(specs=[('NumberOfEmployees', 'cast:string')])
步骤2:禁用Spark向量化读取(若仍报错)
如果Parquet文件中该列存在格式异常值,Spark向量化读取会在转换前触发校验,需在读取时禁用该功能:
datasource0 = glueContext.create_dynamic_frame.from_options( format_options={"enableVectorizedReader": False}, # 禁用向量化读取 connection_type="s3", format="parquet", connection_options={ "paths": s3_path_lst, "recurse": True, }, transformation_ctx="datasource0", )
步骤3:合并Parquet文件Schema(若存在类型不一致)
部分Parquet文件可能存在NumberOfEmployees列类型不一致的情况(如有的是int,有的是long或string),需先合并Schema再转换:
# 开启Schema合并读取所有文件 datasource0 = glueContext.create_dynamic_frame.from_options( format_options={"mergeSchema": True}, connection_type="s3", format="parquet", connection_options={ "paths": s3_path_lst, "recurse": True, }, transformation_ctx="datasource0", ) # 执行类型转换 datasource_dyf = datasource0.resolveChoice(specs=[('NumberOfEmployees', 'cast:string')])
内容的提问来源于stack exchange,提问作者Prathap Veera raghavan
相关产品推荐
相关产品推荐

