You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

将重叠时间区间转换为连续区间:双数据源对账及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;

方案说明

  1. all_time_points:收集两个表中所有的valid_from和valid_to+1日期点,这些点是拆分连续区间的关键节点。
  2. continuous_intervals:通过窗口函数lead将相邻时间点组合成连续区间,确定每个区间的起止范围。
  3. interval_data:将连续区间与source1、source2做左关联,判断区间是否落在表的有效时间范围内,获取对应数据。
  4. 最终查询:通过COALESCE优先取source1的source和value,按时间顺序输出结果,完全匹配需求。

内容的提问来源于stack exchange,提问作者legends1337

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 18:40:00