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
相关产品推荐
相关产品推荐

