将重叠时间区间转换为连续区间:双数据源对账及Hive查询需求
兼容Hive的时间区间对账查询方案
表结构与测试数据
source1表
CREATE TABLE source1 (source string, id string, value string, valid_from STRING, valid_to STRING); INSERT INTO source1 VALUES ('Source1','ABC123', 'IR10', '2020-06-07', '2020-07-07'), ('Source1','ABC123', 'IR29', '2020-07-08', '2020-10-01'), ('Source1','ABC123', 'IR98', '2020-10-02', '9999-12-31');
source2表
CREATE TABLE source2 (source string, id string, value string, valid_from STRING, valid_to STRING); INSERT INTO source2 VALUES ('Source2','ABC123', 'ES23', '2020-04-01', '2020-04-27'), ('Source2','ABC123', 'JD94', '2020-04-28', '2020-05-31'), ('Source2','ABC123', 'JD74', '2020-06-01', '2020-09-29'), ('Source2','ABC123', 'OP13', '2020-09-30', '2020-11-30');
需求说明
对两个数据源进行对账,将重叠的时间区间拆分转换为连续的有效性范围:
- 当source1与source2区间重叠时,以source1为主数据源(最终结果的
source和value取source1的值) - 无重叠的区间保留各自数据源的信息
- 需输出每个连续区间对应的
s1_value(source1的value)、s2_value(source2的value)
预期输出
+---------+--------+-------+----------+----------+------------+------------+ | source | id | value | s1_value | s2_value | valid_from | valid_to | +---------+--------+-------+----------+----------+------------+------------+ | Source2 | ABC123 | ES23 | NULL | ES23 | 2020-04-01 | 2020-04-27 | | Source2 | ABC123 | JD94 | NULL | JD94 | 2020-04-28 | 2020-05-31 | | Source2 | ABC123 | JD74 | NULL | JD74 | 2020-06-01 | 2020-06-06 | | Source1 | ABC123 | IR10 | IR10 | JD74 | 2020-06-07 | 2020-07-07 | | Source1 | ABC123 | IR29 | IR29 | JD74 | 2020-07-08 | 2020-09-29 | | Source1 | ABC123 | IR29 | IR29 | OP13 | 2020-09-30 | 2020-10-01 | | Source1 | ABC123 | IR98 | IR98 | OP13 | 2020-10-02 | 2020-11-30 | | Source1 | ABC123 | IR98 | IR98 | NULL | 2020-12-01 | 9999-12-31 | +---------+--------+-------+----------+----------+------------+------------+
原方案问题
原尝试使用FULL OUTER JOIN编写的查询在Impala中执行报错:
原查询语句
SELECT s1.id AS s1_id, s2.id AS s2_id, COALESCE(s1.source, s2.source) AS source, COALESCE(s1.value, s2.value) AS value, s1.value AS s1_value, s2.value AS s2_value, s1.valid_from AS s1_valid_from, s1.valid_to AS s1_valid_to, s2.valid_from AS s2_valid_from, s2.valid_to AS s2_valid_to, CASE WHEN COALESCE(s1.valid_from, '2020-04-01') < s2.valid_from THEN s2.valid_from ELSE s1.valid_from END AS valid_from, CASE WHEN COALESCE(s1.valid_to, '9999-12-31') > s2.valid_to THEN s2.valid_to ELSE COALESCE(s1.valid_to, '9999-12-31') END AS valid_to FROM source1 s1 FULL OUTER JOIN source2 s2 ON s1.id = s1.id AND ( s1.valid_from BETWEEN s2.valid_from AND s2.valid_to OR s1.valid_to BETWEEN s2.valid_from AND s2.valid_to OR ( s1.valid_from < s2.valid_from AND s1.valid_to > s2.valid_to ) );
报错信息
NotImplementedException: Error generating a valid execution plan for this query. A FULL OUTER JOIN type with no equi-join predicates can only be executed with a single node plan.
兼容Hive的解决方案
思路:先收集所有时间节点生成连续区间,再将每个区间与两个表关联获取对应值,最后按需求处理source1的优先级。
查询语句
WITH all_time_points AS ( -- 收集所有有效起始和结束日期(含结束日期+1,用于生成连续区间) SELECT id, valid_from AS point FROM source1 UNION SELECT id, date_add(valid_to, 1) AS point FROM source1 UNION SELECT id, valid_from AS point FROM source2 UNION SELECT id, date_add(valid_to, 1) AS point FROM source2 ), continuous_intervals AS ( -- 生成连续时间区间 SELECT id, point AS valid_from, lead(point) OVER (PARTITION BY id ORDER BY point) AS valid_to FROM all_time_points WHERE point <= '9999-12-31' ), interval_data AS ( -- 关联两个表,获取每个区间对应的s1和s2数据 SELECT ci.id, ci.valid_from, date_add(ci.valid_to, -1) AS valid_to, s1.source AS s1_source, s1.value AS s1_value, s2.source AS s2_source, s2.value AS s2_value FROM continuous_intervals ci LEFT JOIN source1 s1 ON ci.id = s1.id AND ci.valid_from <= s1.valid_to AND date_add(ci.valid_to, -1) >= s1.valid_from LEFT JOIN source2 s2 ON ci.id = s2.id AND ci.valid_from <= s2.valid_to AND date_add(ci.valid_to, -1) >= s2.valid_from WHERE ci.valid_to IS NOT NULL ) -- 按需求输出结果,优先取source1的信息 SELECT COALESCE(s1_source, s2_source) AS source, id, COALESCE(s1_value, s2_value) AS value, s1_value, s2_value, valid_from, valid_to FROM interval_data ORDER BY valid_from;
方案说明
- all_time_points:收集两个表中所有的
valid_from和valid_to+1日期点,这些点是拆分连续区间的关键节点。 - continuous_intervals:通过窗口函数
lead将相邻时间点组合成连续区间,确定每个区间的起止范围。 - interval_data:将连续区间与source1、source2做左关联,判断区间是否落在表的有效时间范围内,获取对应数据。
- 最终查询:通过
COALESCE优先取source1的source和value,按时间顺序输出结果,完全匹配需求。
内容的提问来源于stack exchange,提问作者legends1337
相关产品推荐
相关产品推荐

