使用PySpark向Hive表写数据时报LongWritable转DoubleWritable错误如何解决
根因说明
报错org.apache.hadoop.io.LongWritable cannot be cast to org.apache.hadoop.io.DoubleWritable的核心原因是Spark写入Hive时,待写入字段的顺序/实际类型和Hive表对应位置声明的类型不匹配:PySpark的insertInto接口按字段位置匹配而非列名匹配,即便你做了列级的类型转换,只要DataFrame和Hive表的列顺序不一致,就会出现类型错位报错,开发环境和测试环境因为DataFrame默认列顺序的差异会出现一边正常一边报错的情况。
可落地解决方案
- 优先替换写入接口,改用按列名匹配的
saveAsTable的append模式,避免列顺序错位问题,代码修改如下:
change_results.select(st_cols).write.mode('append').saveAsTable(args.change_table)
- 若必须使用
insertInto,则显式对齐DataFrame和Hive表的字段顺序,修改类型转换逻辑,直接使用Hive表的列顺序做类型转换而非原DataFrame的列顺序:
from pyspark.sql.types import StringType, LongType, DoubleType import pyspark.sql.functions as F col_map = {'id':StringType(), 'booking_time':StringType(), 'perd':LongType(), 'sequence':LongType(), 'm':StringType(), 'a':DoubleType(), 'b':DoubleType(), 'c':DoubleType(), 'd':DoubleType(), 'e':DoubleType(), 'f':DoubleType(), 'g': DoubleType(), 'h':DoubleType(), 'record_time':StringType()} # 优先获取Hive表的列顺序,以此为基准做类型转换和字段选择 hive_cols = spark.sql(f"select * from {args.change_table} limit 0").columns change_results = change_results.select([F.col(c).cast(col_map[c]).alias(c) for c in hive_cols]) # 写入逻辑保持不变即可 logger.info("Writing") change_results.write.mode('append').insertInto(args.change_table)
- 检查两个环境的Spark参数
spark.sql.storeAssignmentPolicy的差异,测试环境如果配置为LEGACY会触发严格的类型校验,可临时调整为ANSI兼容模式验证:
-- 会话级别调整参数 set spark.sql.storeAssignmentPolicy=ANSI;
- 排查测试环境Hive表是否存在历史脏数据,比如之前写入的Long类型数据存放在Double类型的字段中,导致追加写入时元数据校验失败,可先向测试环境表写入一条全字段符合类型的测试数据验证表本身是否正常。
内容的提问来源于stack exchange,提问作者Ironman
相关产品推荐
相关产品推荐

