Impala中排除分组字段实现Count Distinct统计的方法
问题分析与解决:Impala中COUNT(DISTINCT)结果异常及优化方案
问题描述
现有一段Impala脚本化SQL,通过UNION ALL合并4张表,原查询按原始时间字段等维度分组统计行数。尝试添加COUNT(DISTINCT datasourceid) AS num_datasets后,返回结果中每一行的num_datasets值均为1,疑问是否由GROUP BY中的日期字段导致,同时希望在不将原始插入/更新/源时间字段加入GROUP BY的前提下,正确统计该列的distinct值。
原查询代码
SELECT 'cluster1' AS cluster_name, 'database1' AS database_name, table_name, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name, COUNT(*) AS num_rows FROM ( SELECT 'table1' AS table_name, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table1 UNION ALL SELECT 'table2' AS table_name, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table2 UNION ALL SELECT 'table3' AS table_name, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table3 UNION ALL SELECT 'table4' AS table_name, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table4 ) AS combined_results GROUP BY table_name, etl_update_time, etl_source_time, etl_insert_time, data_provider_location, data_provider_name;
修改后异常代码
SELECT 'cluster1' AS cluster_name, 'database1' AS database_name, table_name, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name, COUNT(*) AS num_rows, COUNT (DISTINCT datasourceid) AS num_datasets, FROM ( SELECT 'table1' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table1 UNION ALL SELECT 'table2' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table2 UNION ALL SELECT 'table3' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table3 UNION ALL SELECT 'table4' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table4 ) AS combined_results GROUP BY table_name, datasourceid, etl_update_time, etl_source_time, etl_insert_time, data_provider_location, data_provider_name;
原因分析
- COUNT(DISTINCT)结果为1的直接原因:你将
datasourceid加入了GROUP BY子句,这会让每个分组内仅包含唯一的datasourceid值,因此COUNT(DISTINCT datasourceid)必然返回1,和日期字段无关。 - 分组粒度冗余问题:原查询中GROUP BY使用的是原始的
etl_insert_time等时间字段,而SELECT输出的是格式化后的小时级时间,这会导致分组粒度远细于实际需求(同一小时内的不同时间戳会被拆分为不同分组)。
解决方案
要实现按小时级时间分组,同时正确统计每个分组内的唯一datasourceid数量,需做以下调整:
- 移除GROUP BY中的
datasourceid和原始时间字段 - 改用格式化后的小时级时间表达式作为分组维度(或在子查询中提前转换时间,简化外层GROUP BY)
优化后代码(方式一:外层直接使用格式化表达式分组)
SELECT 'cluster1' AS cluster_name, 'database1' AS database_name, table_name, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name, COUNT(*) AS num_rows, COUNT(DISTINCT datasourceid) AS num_datasets FROM ( SELECT 'table1' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table1 UNION ALL SELECT 'table2' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table2 UNION ALL SELECT 'table3' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table3 UNION ALL SELECT 'table4' AS table_name, datasourceid, etl_insert_time, etl_update_time, etl_source_time, data_provider_location, data_provider_name FROM test_db.table4 ) AS combined_results GROUP BY table_name, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH'), from_timestamp(etl_update_time, 'yyyy-MM-dd HH'), from_timestamp(etl_source_time, 'yyyy-MM-dd HH'), CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT), CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END, data_provider_name;
优化后代码(方式二:子查询提前转换时间,简化外层分组)
SELECT 'cluster1' AS cluster_name, 'database1' AS database_name, table_name, etl_insert_time_hour, etl_update_time_hour, etl_source_time_hour, time_delay, data_provider_location, data_provider_name, COUNT(*) AS num_rows, COUNT(DISTINCT datasourceid) AS num_datasets FROM ( SELECT 'table1' AS table_name, datasourceid, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name FROM test_db.table1 UNION ALL SELECT 'table2' AS table_name, datasourceid, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name FROM test_db.table2 UNION ALL SELECT 'table3' AS table_name, datasourceid, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name FROM test_db.table3 UNION ALL SELECT 'table4' AS table_name, datasourceid, from_timestamp(etl_insert_time, 'yyyy-MM-dd HH') AS etl_insert_time_hour, from_timestamp(etl_update_time, 'yyyy-MM-dd HH') AS etl_update_time_hour, from_timestamp(etl_source_time, 'yyyy-MM-dd HH') AS etl_source_time_hour, CAST((unix_timestamp(etl_update_time) - unix_timestamp(etl_source_time)) / 3600 AS BIGINT) AS time_delay, CASE WHEN data_provider_location = '' THEN NULL ELSE data_provider_location END AS data_provider_location, data_provider_name FROM test_db.table4 ) AS combined_results GROUP BY table_name, etl_insert_time_hour, etl_update_time_hour, etl_source_time_hour, time_delay, data_provider_location, data_provider_name;
以上两种方式均实现了:
- 按小时级时间维度分组,无需保留原始时间字段在GROUP BY中
- 正确统计每个分组内的唯一
datasourceid数量,不会再返回1
内容的提问来源于stack exchange,提问作者rvnd
相关产品推荐
相关产品推荐

