从DMS导出的Parquet文件创建PySpark DataFrame遇类型兼容问题
AWS DMS导出Parquet兼容PySpark问题解决方案
问题背景
我们通过AWS DMS读取RDS MySQL数据库数据,将其以Parquet格式输出至S3存储桶;随后计划用PySpark读取该文件生成DataFrame,进而创建Hudi数据集实现增量检测,相关代码如下:
%%configure -f { "conf": { "spark.jars":"hdfs:///user/hadoop/aws-java-sdk-bundle-1.12.31.jar, hdfs:///user/hadoop/hudi-spark-bundle.jar,hdfs:///user/hadoop/spark-avro.jar", "spark.sql.hive.convertMetastoreParquet":"false", "spark.serializer":"org.apache.spark.serializer.KryoSerializer", "spark.dynamicAllocation.executorIdleTimeout": 3600, "spark.executor.memory": "5G", "spark.executor.cores": 3, "spark.dynamicAllocation.initialExecutors":5 } } config = { "table_name": "ticket_table", "target": "s3://dms-rds-s3/hudi/hudi_test", "primary_key": "storeid", "sort_key": "ticket_updated_date", "commits_to_retain": "4" } # General Constants HUDI_FORMAT = "org.apache.hudi" TABLE_NAME = "hoodie.table.name" RECORDKEY_FIELD_OPT_KEY = "hoodie.datasource.write.recordkey.field" PRECOMBINE_FIELD_OPT_KEY = "hoodie.datasource.write.precombine.field" OPERATION_OPT_KEY = "hoodie.datasource.write.operation" BULK_INSERT_OPERATION_OPT_VAL = "bulk_insert" UPSERT_OPERATION_OPT_VAL = "upsert" DELETE_OPERATION_OPT_VAL = "delete" BULK_INSERT_PARALLELISM = "hoodie.bulkinsert.shuffle.parallelism" UPSERT_PARALLELISM = "hoodie.upsert.shuffle.parallelism" S3_CONSISTENCY_CHECK = "hoodie.consistency.check.enabled" HUDI_CLEANER_POLICY = "hoodie.cleaner.policy" KEEP_LATEST_COMMITS = "KEEP_LATEST_COMMITS" KEEP_LATEST_FILE_VERSIONS = "KEEP_LATEST_FILE_VERSIONS" HUDI_COMMITS_RETAINED = "hoodie.cleaner.commits.retained" HUDI_FILES_RETAINED = "hoodie.cleaner.fileversions.retained" PAYLOAD_CLASS_OPT_KEY = "hoodie.datasource.write.payload.class.key()" EMPTY_PAYLOAD_CLASS_OPT_VAL = "org.apache.hudi.EmptyHoodieRecordPayload" # Hive Constants HIVE_SYNC_ENABLED_OPT_KEY="hoodie.datasource.hive_sync.enable" HIVE_PARTITION_FIELDS_OPT_KEY="hoodie.datasource.hive_sync.partition_fields" HIVE_ASSUME_DATE_PARTITION_OPT_KEY="hoodie.datasource.hive_sync.assume_date_partitioning" HIVE_PARTITION_EXTRACTOR_CLASS_OPT_KEY="hoodie.datasource.hive_sync.partition_extractor_class" HIVE_TABLE_OPT_KEY="hoodie.datasource.hive_sync.table" # Partition Constants NONPARTITION_EXTRACTOR_CLASS_OPT_VAL="org.apache.hudi.hive.NonPartitionedExtractor" MULTIPART_KEYS_EXTRACTOR_CLASS_OPT_VAL="org.apache.hudi.hive.MultiPartKeysValueExtractor" KEYGENERATOR_CLASS_OPT_KEY="hoodie.datasource.write.keygenerator.class" NONPARTITIONED_KEYGENERATOR_CLASS_OPT_VAL="org.apache.hudi.keygen.NonpartitionedKeyGenerator" COMPLEX_KEYGENERATOR_CLASS_OPT_VAL="org.apache.hudi.ComplexKeyGenerator" PARTITIONPATH_FIELD_OPT_KEY="hoodie.datasource.write.partitionpath.field" #Incremental Constants VIEW_TYPE_OPT_KEY="hoodie.datasource.query.type" BEGIN_INSTANTTIME_OPT_KEY="hoodie.datasource.read.begin.instanttime" VIEW_TYPE_INCREMENTAL_OPT_VAL="incremental" END_INSTANTTIME_OPT_KEY="hoodie.datasource.read.end.instanttime" df1 = sqlContext.read.parquet("PATH")
错误信息
读取S3上的Parquet文件时触发如下错误:
"An error was encountered: An error occurred while calling o90.parquet. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.3 in stage 0.0 (TID 3) (ip-10-240-84-185.us-west-2.compute.internal executor 1): org.apache.spark.sql.AnalysisException: Parquet type not supported: INT32 (UINT_32)"
原因分析
- 核心问题是类型映射不兼容:AWS DMS默认会将MySQL中的无符号整数类型(如
UNSIGNED INT)转换为Parquet的UINT_32类型,但Spark的Parquet解析器原生不支持无符号整数类型,导致解析时抛出类型不支持的异常。
解决方案
方案1:调整AWS DMS设置(推荐)
从源头修改Parquet输出的类型映射规则,让DMS生成Spark兼容的格式:
- 登录AWS DMS控制台,找到对应的迁移任务,进入目标端点的编辑页面。
- 在Extra connection attributes中添加配置项:
parquetEnableUInt32=false。该配置会强制DMS将MySQL无符号整数类型映射为Parquet的有符号INT32类型,完全兼容Spark的解析逻辑。 - 保存配置后重启DMS任务,重新生成Parquet文件即可正常读取。
方案2:Spark端兼容处理(备选)
如果无法修改DMS任务配置,可在Spark读取时手动处理类型兼容:
- 自定义Schema读取:
手动定义DataFrame的Schema,将原UINT_32类型的字段指定为Spark支持的IntegerType或LongType,示例代码如下:from pyspark.sql.types import StructType, StructField, IntegerType, TimestampType, StringType # 根据实际表结构定义Schema,替换UINT_32字段为IntegerType custom_schema = StructType([ StructField("storeid", IntegerType(), nullable=False), StructField("ticket_updated_date", TimestampType(), nullable=True), # 添加其他字段定义 ]) df1 = sqlContext.read.schema(custom_schema).parquet("s3://dms-rds-s3/your-parquet-path/") - 禁用向量化读取:
在Spark配置中添加spark.sql.parquet.enableVectorizedReader=false,关闭向量化读取以兼容无符号整数类型,但此方法会降低读取性能,仅适合小数据量场景。
内容的提问来源于stack exchange,提问作者SWATHI S
相关产品推荐
相关产品推荐

