如何编写高效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函数实现:
实现步骤说明
- 生成0-23小时的完整数据集:确保每个小时都有记录,哪怕该小时没有数据
- 预处理原始数据:提取时间小时、去重目标字段、将需要聚合的字段转为结构体
- 左连接+聚合:把原始数据和完整小时数据集关联,聚合每个小时的去重数据为数组
- 生成目标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
相关产品推荐
相关产品推荐

