You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

从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读取时手动处理类型兼容:

  1. 自定义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/")
    
  2. 禁用向量化读取:
    在Spark配置中添加spark.sql.parquet.enableVectorizedReader=false,关闭向量化读取以兼容无符号整数类型,但此方法会降低读取性能,仅适合小数据量场景。

内容的提问来源于stack exchange,提问作者SWATHI S

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 18:48:38