Spark SQL大数据聚合查询优化:时间戳错误与性能难题
微秒级时间戳Spark SQL高效聚合方案问题
我正尝试用Spark SQL从Azure数据库聚合微秒级时间戳的数据,遇到以下问题:
初始查询:逻辑目标正确但结果错误
初始参考查询逻辑是按指定时间间隔(由整数乘数N和时间单位秒数U计算,如分钟=60、小时=3600)分组聚合,但本地转换时间戳时结果不正确,无法定位具体逻辑问题:
select max(Timestamp), mean(datacolumn1), median(datacolumn2), max(datacolumn3) from myTable where datacolumn2 = "Yes" and Timestamp between "1684738800000000" and "1684825200000000" GROUP BY FLOOR((Timestamp)/(N*U*1E6))*(N*U*1E6)
重写查询:结果正确但性能极差
完全重写后的查询结果准确,但性能比初始查询慢约120倍,处理数十亿行月度/年度数据时效率无法接受:
WITH Aggregated AS ( -- Filter data, convert timestamps, and aggregate in one step SELECT FLOOR(Timestamp / (1000000 * N * U)) * (N * U) AS grouped_timestamp, min(Timestamp) as min_timestamp, AVG(data1) AS data1, AVG(data2) AS data2, AVG(data3) AS data3 FROM dataBaseTable WHERE Timestamp BETWEEN UNIXtstart AND UNIXtstop GROUP BY FLOOR(Timestamp / (1000000 * N * U)) ) SELECT grouped_timestamp, from_unixtime(CAST(grouped_timestamp / 1000000 AS BIGINT)) AS datetime_representation, data1, data2, data3 FROM Aggregated ORDER BY grouped_timestamp;
优化查询:速度快但时间戳计算错误
尝试优化后的查询速度达标,但生成错误的时间戳(如-8741307600),推测是整数溢出问题,即使显式将Timestamp转为BIGINT也无法解决:
WITH Aggregated AS ( -- Filter data, convert timestamps, and aggregate in one step SELECT CAST(FLOOR(CAST(Timestamp AS BIGINT) / (1000000 * N * U)) * (N * U) AS BIGINT) AS grouped_timestamp, MIN(Timestamp) as min_timestamp, AVG(data1) AS data1, AVG(data2) AS data2, AVG(data3) AS data3 FROM dataBaseTable WHERE Timestamp BETWEEN UNIXtstart AND UNIXtstop GROUP BY FLOOR(CAST(Timestamp AS BIGINT) / (1000000 * N * U)) ) SELECT grouped_timestamp, from_unixtime(CAST(grouped_timestamp / 1000000 AS BIGINT)) AS datetime_representation, data1, data2, data3 FROM Aggregated ORDER BY grouped_timestamp;
需求
需要实现针对数十亿行微秒级数据集的准确且高效的聚合方案,其中UNIXtstart、UNIXtstop、N和U均为用户输入的可变参数。
内容的提问来源于stack exchange,提问作者Sophia
相关产品推荐
相关产品推荐

