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

Spark-JAVA中基于两个Dataset实现指定SQL查询的数据转换

Spark Java实现SQL转换逻辑的方案

没问题,我来帮你把这段SQL转换成Spark Java的Dataset API操作,步骤很清晰,咱们一步步来:

第一步:处理子查询获取最大日期

首先要拿到record_status表中db_destination_name='AB'的最大date_record,如果没有匹配记录就用默认值'0000-00-00 00:00:00'。用Spark的聚合函数就能轻松实现:

import static org.apache.spark.sql.functions.*;

// 过滤目标记录并聚合得到最大日期(含默认值)
Dataset<Row> maxDateDataset = ds_record_status
    .filter(col("db_destination_name").equalTo("AB"))
    .agg(coalesce(max(col("date_record")), lit("0000-00-00 00:00:00")).alias("max_date"));

// 提取单个日期值,方便后续过滤
String targetMaxDate = maxDateDataset.first().getString(0);

第二步:主数据集的过滤与转换

接下来对ds_evenement做过滤、列重命名和空值替换,完全对应SQL里的逻辑:

// 构建最终结果数据集
Dataset<Row> resultDataset = ds_evenement
    // 过滤DATE_HIST大于目标日期的记录
    .filter(col("DATE_HIST").gt(lit(targetMaxDate)))
    // 列处理:别名+空值替换为空字符串
    .select(
        col("ID").alias("Identifier"),
        coalesce(col("INTITULE"), lit("")).alias("NAME_INTITULE"),
        coalesce(col("ID_CAT"), lit("")).alias("CODE_CATEGORIE")
    );

可选优化:用Join替代提取单值

如果不想把日期值提取出来(比如担心并发场景下数据变化),可以用Left Semi Join结合广播优化来实现过滤,性能也很不错:

// 先构建子查询数据集
Dataset<Row> subQuery = ds_record_status
    .filter(col("db_destination_name").equalTo("AB"))
    .agg(coalesce(max(col("date_record")), lit("0000-00-00 00:00:00")).alias("max_date"));

// 通过Left Semi Join过滤主表,同时广播小数据集提升性能
Dataset<Row> resultDataset = ds_evenement.join(
    broadcast(subQuery),
    ds_evenement.col("DATE_HIST").gt(subQuery.col("max_date")),
    "left_semi"
)
.select(
    col("ID").alias("Identifier"),
    coalesce(col("INTITULE"), lit("")).alias("NAME_INTITULE"),
    coalesce(col("ID_CAT"), lit("")).alias("CODE_CATEGORIE")
);

关键说明

  • Spark里的coalesce函数和SQL的IFNULL逻辑完全一致,都是返回第一个非空值;
  • Left Semi Join只会保留主表中符合关联条件的记录,不会引入多余字段,非常适合这种过滤场景;
  • 用broadcast广播小的子查询数据集,可以避免shuffle,大幅提升性能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 06:31:26