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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 15:37:14