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

如何编写高效Spark SQL查询实现按小时聚合去重的JSON结果

高效实现Spark SQL按小时分组去重并生成目标JSON结构

我有一个结构为[{"time","currentStop","lat","lon","speed"}]的JSON文件,示例数据如下:

[ {"time":"2015-06-09 23:59:59","currentStop":"xx","lat":"22.264856","lon":"113.520450","speed":"25.30"}, {"time":"2015-06-09 21:00:49","currentStop":"yy","lat":"22.263","lon":"113.52","speed":"34.5"}, {"time":"2015-06-09 21:55:49","currentStop":"zz","lat":"21.3","lon":"113.521","speed":"13.7"} ]

期望得到结构为[{"hour","value":["currentStop","lat","lon","speed"]}]的JSON结果,展示每小时的(currentStop,lat,lon,speed)去重数据,示例结果如下:

[ {"hour":0,"value":[]}, {"hour":1,"value":[]}, ...... {"hour":21,"value":[{"currentStop":"yy","lat":"22.263","lon":"113.52","speed":"34.5"},{"currentStop":"zz","lat":"21.3","lon":"113.521","speed":"13.7"}]} {"hour":23, "value": [{"currentStop":"xx","lat":22.264856,"lon":113.520450,"speed":25.30}]}, ]

我使用Spark Java API,通过循环方式已实现但效率极低、成本高,现有代码如下:

Dataset<Row> bus_ic=spark.read().json(file); 
bus_ic.createOrReplaceTempView("view"); 
StringBuilder text = new StringBuilder("["); 
bus_ic.select(bus_ic.col("currentStop"), bus_ic.col("lon").cast("double"), bus_ic.col("speed").cast("double"), bus_ic.col("lat").cast("double"),bus_ic.col("LINEID"), bus_ic.col("time").cast("timestamp")) 
      .createOrReplaceTempView("view"); 
StringBuilder sqlString = new StringBuilder(); 
for(int i = 0; i<24; i++){ 
    sqlString.delete(0,sqlString.length()); 
    sqlString.append("select currentStop, speed, lat, lon from view where hour(time) = ") 
             .append(i) 
             .append(" group by currentStop, speed, lat, lon"); 
    Dataset<Row> t = spark.sql(sqlString.toString()); 
    text.append("{") 
        .append("\"h\":").append(i) 
        .append(",\"value\":") 
        .append(t.toJSON().collectAsList().toString()) 
        .append("}"); 
    if(i!=23) text.append(","); 
} 
text.append("]");

请问如何编写高效的Spark SQL查询实现该目标?


当然可以!你当前的循环方法效率低的核心原因是:重复发起24次Spark查询,每次都要重新扫描全量数据,还频繁把分布式数据拉到Driver端,带来大量冗余开销。下面是一种高效的方案,用一次分布式计算就能完成所有需求,结合Spark DataFrame API和SQL函数实现:

实现步骤说明

  1. 生成0-23小时的完整数据集:确保每个小时都有记录,哪怕该小时没有数据
  2. 预处理原始数据:提取时间小时、去重目标字段、将需要聚合的字段转为结构体
  3. 左连接+聚合:把原始数据和完整小时数据集关联,聚合每个小时的去重数据为数组
  4. 生成目标JSON结构:最后一次性导出结果,避免多次拉取数据到Driver

Java代码实现

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.functions.*;
import java.util.List;
import java.util.stream.Collectors;
import java.util.stream.IntStream;

// 1. 读取原始JSON数据
Dataset<Row> bus_ic = spark.read().json(file);

// 2. 生成0到23的完整小时数据集
List<Integer> hoursList = IntStream.rangeClosed(0, 23).boxed().collect(Collectors.toList());
Dataset<Row> allHours = spark.createDataFrame(hoursList, Integer.class).toDF("hour");

// 3. 预处理原始数据:提取小时、去重、转为结构体
Dataset<Row> processedData = bus_ic
    // 转换time为timestamp类型,提取小时
    .withColumn("time", col("time").cast("timestamp"))
    .withColumn("hour", hour(col("time")))
    // 选择需要的字段,这里可以根据需求转换lat/lon/speed的类型
    .select(
        col("hour"),
        col("currentStop"),
        col("lat").cast("double"), // 可选:如果需要数字类型就转换
        col("lon").cast("double"),
        col("speed").cast("double")
    )
    // 对每个小时内的currentStop+lat+lon+speed去重
    .distinct()
    // 将需要聚合的字段打包成结构体,方便后续聚合为数组
    .withColumn("value_struct", struct(
        col("currentStop"),
        col("lat"),
        col("lon"),
        col("speed")
    ));

// 4. 左连接完整小时数据集,聚合每个小时的value数组
Dataset<Row> result = allHours
    // 左连接确保所有小时都被保留,哪怕没有对应数据
    .join(processedData, allHours.col("hour").equalTo(processedData.col("hour")), "left")
    // 按小时分组,聚合结构体为数组
    .groupBy(allHours.col("hour"))
    .agg(collect_list("value_struct").alias("value"))
    // 把没有数据的小时的value字段转为空数组(替代null)
    .withColumn("value", when(col("value").isNull(), array()).otherwise(col("value")));

// 5. 生成目标JSON字符串
String finalJson = "[" + String.join(",", result.toJSON().collectAsList()) + "]";

效率优势对比

  • 减少作业开销:从24次Spark作业变成1次,避免了重复的任务调度和数据扫描
  • 分布式计算:所有去重、聚合操作都在集群分布式执行,只有最后一步把结果拉到Driver
  • 避免冗余操作:用distinct和collect_list高效完成去重和数组聚合,比循环查询的效率提升数倍(数据量越大,提升越明显)

可选调整点

  • 如果不需要转换lat/lon/speed为数字类型,去掉.cast("double")即可
  • 如果需要保留LINEID字段,把它加入到struct的参数中即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:08:07