Spark Dataset字段插入Hive表顺序错误的问题排查与解决
问题描述
环境:
- Hive 2.1.1-cdh6.2.1
- Spark依赖:
spark-sql_2.11:2.4.0-cdh6.2.1、spark-hive_2.11:2.4.0-cdh6.2.1
Hive建表语句
CREATE TABLE COUNT_REPORT ( PROCESSING_DAY INT, TOTAL_UIT_IDS INT, DISTINCT_UIT_IDS INT, DISTINCT_UIT_IDS_TRADE_AGREEMENT_RELATION_TOTAL INT, DISTINCT_UIT_IDS_TRADE_AGREEMENT_RELATION_FOR_PROCESSING_DAY INT ) STORED AS PARQUET;
Java写入代码
DTO类:
@Data @AllArgsConstructor public class CountReportItem implements Serializable { private long processingDate; private long totalOptimaUitIds; private long distinctOptimaUitIds; private long countOfDistinctUitIdsInTradeAgreements; private long countOfDistinctUitIdsInTradeAgreementsForTradeDate; }
数据写入逻辑:
long processingDate = /*...*/; long totalUitIds = /*...*/; long distinctUitIds = /*...*/; long countOfDistinctUitIdsInTradeAgreements = /*...*/; long countOfDistinctUitIdsInTradeAgreementsForProcDate = /*...*/; CountReportItem countReportItem = new CountReportItem( processingDate, totalUitIds, distinctUitIds, countOfDistinctUitIdsInTradeAgreements, countOfDistinctUitIdsInTradeAgreementsForProcDate ); sparkSession.createDataset(Arrays.asList(countReportItem), Encoders.bean(CountReportItem.class)) .write() .format("parquet") .option("compression", "snappy") .mode(SaveMode.Append) .insertInto("COUNT_REPORT");
问题现象
执行Hive查询:
hive> set hive.cli.print.header=true; hive> select * from COUNT_REPORT;
得到结果:
processing_day total_uit_ids distinct_uit_ids distinct_uit_ids_trade_agreement_relation_total distinct_uit_ids_trade_agreement_relation_for_processing_day 1372171 1372171 826053 20221230 1496602 1436195 1436195 870445 20230227 1574622
其中20221230、20230227是PROCESSING_DAY的正确值,却被写入了第四列,说明DTO字段按字母顺序而非声明顺序插入Hive表,导致数据错位。
核心原因
Spark 的 insertInto 方法是按列的位置顺序匹配Hive表的列,而非根据列名映射。而Encoders.bean(CountReportItem.class)生成Dataset时,默认通过反射按字母顺序获取Java类的字段,导致Dataset的列顺序与Hive表的列顺序不一致,最终数据错位。
解决方法
方法一:手动构建DataFrame并指定列顺序
放弃Encoders.bean,手动定义与Hive表列顺序完全一致的Schema,将数据转换为Row后创建DataFrame:
import org.apache.spark.sql.Row; import org.apache.spark.sql.RowFactory; import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.types.StructType; import java.util.Arrays; // 定义与Hive表列顺序一致的Schema StructType schema = new StructType() .add("processing_day", DataTypes.LongType) .add("total_uit_ids", DataTypes.LongType) .add("distinct_uit_ids", DataTypes.LongType) .add("distinct_uit_ids_trade_agreement_relation_total", DataTypes.LongType) .add("distinct_uit_ids_trade_agreement_relation_for_processing_day", DataTypes.LongType); // 将数据转换为Row Row row = RowFactory.create( processingDate, totalUitIds, distinctUitIds, countOfDistinctUitIdsInTradeAgreements, countOfDistinctUitIdsInTradeAgreementsForProcDate ); // 创建DataFrame并插入 sparkSession.createDataFrame(Arrays.asList(row), schema) .write() .format("parquet") .option("compression", "snappy") .mode(SaveMode.Append) .insertInto("COUNT_REPORT");
方法二:使用@JsonPropertyOrder指定DTO序列化顺序
如果要保留Encoders.bean,可以添加Jackson的@JsonPropertyOrder注解,强制指定DTO字段的序列化顺序,使其与Hive表列顺序对应:
import com.fasterxml.jackson.annotation.JsonPropertyOrder; import lombok.AllArgsConstructor; import lombok.Data; import java.io.Serializable; @Data @AllArgsConstructor @JsonPropertyOrder({ "processingDate", "totalOptimaUitIds", "distinctOptimaUitIds", "countOfDistinctUitIdsInTradeAgreements", "countOfDistinctUitIdsInTradeAgreementsForTradeDate" }) public class CountReportItem implements Serializable { private long processingDate; private long totalOptimaUitIds; private long distinctOptimaUitIds; private long countOfDistinctUitIdsInTradeAgreements; private long countOfDistinctUitIdsInTradeAgreementsForTradeDate; }
添加注解后,Encoders.bean生成的Dataset列顺序会与注解指定的顺序一致,插入时数据即可正确匹配Hive表列位置。
方法三:手动调整Dataset列顺序
在生成Dataset后,通过select方法显式指定列的顺序和别名,确保与Hive表列完全匹配:
import static org.apache.spark.sql.functions.col; sparkSession.createDataset(Arrays.asList(countReportItem), Encoders.bean(CountReportItem.class)) // 按Hive表列顺序映射字段(驼峰转下划线) .select( col("processingDate").alias("processing_day"), col("totalOptimaUitIds").alias("total_uit_ids"), col("distinctOptimaUitIds").alias("distinct_uit_ids"), col("countOfDistinctUitIdsInTradeAgreements").alias("distinct_uit_ids_trade_agreement_relation_total"), col("countOfDistinctUitIdsInTradeAgreementsForTradeDate").alias("distinct_uit_ids_trade_agreement_relation_for_processing_day") ) .write() .format("parquet") .option("compression", "snappy") .mode(SaveMode.Append) .insertInto("COUNT_REPORT");
这种方式既保证了字段名的正确映射,也强制了列顺序与Hive表一致,避免数据错位。
内容的提问来源于stack exchange,提问作者Sergey Tsypanov
相关产品推荐
相关产品推荐

