多车辆行程计数SQL查询异常排查(AWS Athena)
多车辆行程统计的Gaps-and-Islands问题排查与修复
问题背景
需要统计固定时间段内多辆车的有效行程数:车辆Trip_signal_value从非5/6跳转为5/6视为启动,从5/6回落至低于5视为关闭,结合speed字段(行程内有移动即speed总和>0、记录数>1)判断有效行程。单辆车统计正常,但多辆车同时查询时,新增车辆会导致计数结果异常,需排查根本原因(拒绝逐个查询的临时方案)。当前使用AWS Athena(基于Presto/Trino)。
现有SQL语句
WITH cte_data AS ( SELECT car_id, time_stamp, CAST(json_extract_scalar(json_load,'$.Trip_siganl_value') AS INTEGER) AS current_value, LAG(CAST(json_extract_scalar("json_load",'$.Trip_siganl_value') AS INTEGER)) OVER (PARTITION BY car_id ORDER BY time_stamp) as Previousvalue, LEAD(CAST(json_extract_scalar("json_load",'$.Trip_siganl_value') AS INTEGER)) OVER (PARTITION BY car_id ORDER BY time_stamp) as Nextvalue, CAST(json_extract_scalar("json_load",'$.speed') AS DOUBLE) AS speed FROM data_base WHERE "car_id" in ('car_id1', 'car_id2', 'car_id3') ORDER BY time_stamp ), cte_start AS ( SELECT car_id, "time_stamp" as starting_time_location, current_value , ROW_NUMBER() OVER (PARTITION BY car_id ORDER BY time_stamp) AS island_number FROM cte_data WHERE (Previousvalue < 5 OR Previousvalue IS NULL) AND current_value in (5,6) ), cte_end AS ( SELECT car_id, "time_stamp" as ending_time_location, current_value , ROW_NUMBER() OVER ( PARTITION BY car_id ORDER BY time_stamp) AS island_number FROM cte_data WHERE (Nextvalue < 5 OR Nextvalue IS NULL) AND current_value in (5,6) ), cte_final AS ( SELECT cte_start.car_id AS car_id, cte_start.starting_time_location AS start_date_final, cte_end.ending_time_location as end_date_final, ( SELECT COUNT(*) FROM cte_data WHERE cte_data.time_stamp >= cte_start.starting_time_location AND cte_data.time_stamp <= cte_end.ending_time_location ) AS island_row_count, ( SELECT SUM(speed) FROM cte_data WHERE cte_data.time_stamp >= cte_start.starting_time_location AND cte_data.time_stamp <= cte_end.ending_time_location ) AS island_speed_sum FROM cte_start INNER JOIN cte_end ON cte_start.island_number = cte_end.island_number ) SELECT car_id, COUNT(*) FROM cte_final WHERE island_row_count > 1 AND island_speed_sum > 0 GROUP BY car_id
测试数据
car_id time_stamp Trip_signal_value speed ---------------------------------------------- car_id1 1 3 0 car_id1 2 4 0 car_id1 3 5 0 car_id1 4 5 3 car_id1 5 5 5 car_id1 6 5 8 car_id1 7 6 10 car_id1 8 6 14 car_id1 9 5 10 car_id1 10 5 5 car_id1 11 6 3 car_id1 12 6 0 car_id1 13 3 0 car_id1 14 3 0 car_id1 15 3 0 car_id1 16 3 0 car_id1 17 0 0 car_id1 18 0 0 car_id1 19 0 0 car_id1 20 2 0 car_id1 21 4 0 car_id1 22 5 0 car_id1 23 5 3 car_id1 24 5 6 car_id1 25 5 5 car_id1 26 5 9 car_id1 27 5 10 car_id1 28 5 5 car_id1 29 5 4 car_id1 30 5 0 car_id1 31 2 0 car_id1 32 1 0 car_id1 33 0 0 car_id2 7 2 0 car_id2 8 2 0 car_id2 9 5 0 car_id2 10 5 2 car_id2 11 6 4 car_id2 12 6 9 car_id2 13 5 12 car_id2 14 5 13 car_id2 15 5 10 car_id2 16 5 7 car_id2 17 6 4 car_id2 18 5 0 car_id2 19 3 0 car_id2 20 2 0 car_id2 21 2 0 car_id2 22 2 0
期望输出
car_id Number_Of_trips --------------------------- car_id1 2 car_id2 1
问题根因分析
- 子查询未关联
car_id:在cte_final的两个子查询(统计island_row_count和island_speed_sum)中,仅通过时间范围过滤,未限定cte_data.car_id = cte_start.car_id。这会导致统计的是所有车辆在该时间范围内的记录数和speed总和,而非当前车辆的。新增车辆时,时间范围重叠的记录会被计入,直接导致结果异常。 cte_data中的ORDER BY无意义:CTE中的ORDER BY不会保留后续查询的排序顺序,反而可能引发不必要的计算,可直接移除。
修正后的SQL语句
WITH cte_data AS ( SELECT car_id, time_stamp, CAST(json_extract_scalar(json_load,'$.Trip_siganl_value') AS INTEGER) AS current_value, LAG(CAST(json_extract_scalar("json_load",'$.Trip_siganl_value') AS INTEGER)) OVER (PARTITION BY car_id ORDER BY time_stamp) as Previousvalue, LEAD(CAST(json_extract_scalar("json_load",'$.Trip_siganl_value') AS INTEGER)) OVER (PARTITION BY car_id ORDER BY time_stamp) as Nextvalue, CAST(json_extract_scalar("json_load",'$.speed') AS DOUBLE) AS speed FROM data_base WHERE "car_id" in ('car_id1', 'car_id2', 'car_id3') ), cte_start AS ( SELECT car_id, "time_stamp" as starting_time_location, current_value , ROW_NUMBER() OVER (PARTITION BY car_id ORDER BY time_stamp) AS island_number FROM cte_data WHERE (Previousvalue < 5 OR Previousvalue IS NULL) AND current_value in (5,6) ), cte_end AS ( SELECT car_id, "time_stamp" as ending_time_location, current_value , ROW_NUMBER() OVER (PARTITION BY car_id ORDER BY time_stamp) AS island_number FROM cte_data WHERE (Nextvalue < 5 OR Nextvalue IS NULL) AND current_value in (5,6) ), cte_final AS ( SELECT s.car_id, s.starting_time_location, e.ending_time_location, COUNT(d.time_stamp) AS island_row_count, SUM(d.speed) AS island_speed_sum FROM cte_start s JOIN cte_end e ON s.car_id = e.car_id AND s.island_number = e.island_number JOIN cte_data d ON d.car_id = s.car_id AND d.time_stamp BETWEEN s.starting_time_location AND e.ending_time_location GROUP BY s.car_id, s.starting_time_location, e.ending_time_location ) SELECT car_id, COUNT(*) AS Number_Of_trips FROM cte_final WHERE island_row_count > 1 AND island_speed_sum > 0 GROUP BY car_id
修正说明
- 将
cte_final中的关联子查询改为JOIN方式,明确关联car_id,确保只统计当前车辆的行程数据。 - 移除
cte_data中无效的ORDER BY。 - 用
BETWEEN简化时间范围判断,逻辑保持一致。 - 最终输出字段命名与期望结果对齐。
内容的提问来源于stack exchange,提问作者SomeOne
相关产品推荐
相关产品推荐

