欢迎光临
我们一直在努力

Flink 系列第23篇:Flink SQL 多流 Join 指南 —— 从原理到调优,一文吃透

一、Flink Join 简介

在实时数仓的建设过程中,Join 操作无处不在。无论是日志关联扩充维度数据构建宽表,还是通过 ID 关联计算转化率,Join 都是最基础、最核心的操作之一。然而,流处理中的 Join 与离线批处理有着本质区别——数据是无限、持续到达的,且可能随时发生更新或撤回,这使得流 Join 的设计远比离线 Join 复杂。

Flink SQL 支持对动态表进行复杂而灵活的连接操作。为了应对不同场景的语义需求,Flink 提供了多种不同类型的 Join。

从整体来看,Flink 主要支持以下三种 Join 方式:

  • 动态表(流)Join 动态表(流):两条实时数据流之间的关联
  • 动态表(流)Join 外部维表:流数据实时查询外部存储系统(如 Redis、MySQL)获取维度信息
  • 动态表字段的列转行:一种特殊的 Join 形式,将单行中的数组或嵌套结构展开为多行

进一步细分,Flink SQL 支持的 Join 类型包括:

Join 类型分类说明
Regular Join 流与流的 Join 最通用的 Join,支持 INNER/LEFT/RIGHT/FULL OUTER
Interval Join 流与流的 Join 两条流在指定时间区间内的 Join
Temporal Join 流与流的 Join 基于事件时间或处理时间的时态表关联
Lookup Join 流与外部维表的 Join 流数据实时查询外部维表
Array Expansion 列转行 表字段的列转行,类似 Hive 的 explode
Table Function 列转行 自定义函数的列转行,支持 Inner Join 和 Left Outer Join

默认情况下,Join 的顺序没有被优化。表的 Join 顺序是在 FROM 子句中指定的。可以通过把更新频率最低的表放在最前面、频率最高的放在最后这种方式来微调 Join 查询的性能。需要确保表的顺序不会产生笛卡尔积,因为不支持这样的操作并且会导致查询失败。

二、Regular Join(常规 Join)

2.1 定义与原理

Regular Join 是最通用的 Join 类型。在这种 Join 下,Join 两侧表的任何新记录或变更都是可见的,并会影响整个 Join 的结果。例如:如果左边有一条新记录,在关联字段相等的情况下,它将和右边表之前和之后的所有记录进行 Join。

对于流式查询,Regular Join 的语法是最灵活的,允许任何类型的更新(插入、更新、删除)输入表。然而,这种操作具有重要的操作含义:Flink 需要将 Join 输入的两边数据永远保持在状态中。因此,计算查询结果所需的状态可能会无限增长,这取决于所有输入表的输入数据量。

Regular Join 支持四种连接类型:

  • INNER JOIN:只返回两表中满足连接条件的记录(交集)
  • LEFT JOIN:返回左表中的所有记录,即使右表中没有匹配
  • RIGHT JOIN:返回右表中的所有记录,即使左表中没有匹配
  • FULL OUTER JOIN:返回两表的并集,包含匹配和不匹配的记录

Flink 目前只支持等值连接(equi-join),即至少有一个等值条件。不支持任意的 cross join 和 theta join。

2.2 Inner Join 行为

流任务中,只有两条流 Join 到才输出,输出 +[L, R]。

CREATE TABLE show_log_table (
log_id BIGINT,
show_params STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '1',
'fields.show_params.length' = '3',
'fields.log_id.min' = '1',
'fields.log_id.max' = '10'
);

CREATE TABLE click_log_table (
log_id BIGINT,
click_params STRING
) WITH (
'connector' = 'datagen',
'rows-per-second' = '1',
'fields.click_params.length' = '3',
'fields.log_id.min' = '1',
'fields.log_id.max' = '10'
);

CREATE TABLE sink_table (
s_id BIGINT,
s_params STRING,
c_id BIGINT,
c_params STRING
) WITH (
'connector' = 'print'
);

INSERT INTO sink_table
SELECT
show_log_table.log_id AS s_id,
show_log_table.show_params AS s_params,
click_log_table.log_id AS c_id,
click_log_table.click_params AS c_params
FROM show_log_table
INNER JOIN click_log_table
ON show_log_table.log_id = click_log_table.log_id;

2.3 Left Join 与回撤流

Left Join 的执行逻辑更为复杂:

  • 左流数据到达,无论有没有 Join 到右流的数据,都会输出:
    • Join 到:输出 +[L, R]
    • 没 Join 到:输出 +[L, null]
  • 如果之后右流数据到达,发现左流之前输出过没有 Join 到的数据,则会发起回撤流:先输出 -[L, null],然后输出 +[L, R]

回撤流(Retract Stream)是 Flink 流处理中一个非常重要的概念。它通过发送“撤回”标记来撤销之前已经输出的结果,然后再输出更正后的新结果。这种机制保证了流式计算结果最终能与批处理结果一致。

RIGHT JOIN 与 LEFT JOIN 逻辑一致,只是左表和右表的角色互换。

2.4 Full Join 行为

Full Outer Join 最为复杂。左流或右流的数据到达之后,无论有没有 Join 到另外一条流的数据,都会输出:

  • 对右流来说:Join 到输出 +[L, R],没 Join 到输出 +[null, R]
  • 对左流来说:Join 到输出 +[L, R],没 Join 到输出 +[L, null]

如果一条流的数据到达之后,发现之前另一条流输出过没有 Join 到的数据,则会发起回撤流:

  • 左流数据到达:回撤 -[null, R],输出 +[L, R]
  • 右流数据到达:回撤 -[L, null],输出 +[L, R]

2.5 Regular Join 的注意事项

  • 会产生回撤流:外连接(LEFT/RIGHT/FULL)可能产生回撤流,下游系统需要支持回撤处理。

  • 等值 Join vs 非等值 Join:实时 Regular Join 可以不是等值 Join。其区别在于:

    • 等值 Join:数据 Shuffle 策略是 Hash,按照等值条件中的关联键发往对应的下游算子
    • 非等值 Join:数据 Shuffle 策略是 Global,所有数据发往一个并发,按照非等值条件进行关联
  • Join 的流程:左流新来一条数据之后,会和右流中符合条件的所有数据做 Join,然后输出。

  • State 无限增长风险:流的上游是无限的数据。Flink 会将两条流的所有数据都存储在 State 中,所以 State 会无限增大。因此需要为 State 配置合适的 TTL(Time-To-Live),以防止 State 过大。但需要注意:设置 TTL 可能会影响查询结果的正确性。

  • Join 顺序优化:在 FROM 子句中,将更新频率最低的表放在前面,频率最高的放在后面,有助于提升性能。

  • 三、Interval Join(时间区间 Join)

    3.1 定义与原理

    Interval Join 返回一个符合 Join 条件和时间限制的简单笛卡尔积。与 Regular Join 相比,它最大的优势在于可以避免回撤流。

    Interval Join 需要至少一个等值 Join 条件和一个 Join 两边都包含的时间限定 Join 条件。时间范围的判断可以定义为一个条件(如 <, <=, >=, >),也可以使用 BETWEEN 条件,或者两边表的相同类型的时间属性(处理时间或事件时间)的等式判断。

    与 Regular Join 操作相比,Interval Join 只支持带有时间属性的追加表(append-only tables)。由于时间属性是单调递增的,Flink 可以从状态中移除过期的数据,而不会影响结果的正确性。这正是 Interval Join 能够避免 State 无限膨胀的核心原因。

    3.2 执行行为详解

    以左流 L、右流 R 为例:

    Inner Interval Join:只有两条流 Join 到时才输出 +[L, R]。需要同时满足 Join 条件中的时间区间条件和其他等值条件。

    Left Interval Join:

    • 左流数据到达,如果没有 Join 到右流数据,会等待(放在 State 中)
    • 如果之后右流数据到达且能 Join 到,则输出 +[L, R]
    • 事件时间模式下,随着 Watermark 推进,如果左流 State 中的数据过期,则删除数据并输出 +[L, null]
    • 如果右流 State 中的数据过期,直接删除

    Right Interval Join:与 Left Interval Join 逻辑相同,但左表和右表的角色互换。

    Full Interval Join:

    • 左流或右流数据到达,如果没有 Join 到另一条流的数据,会等待
    • 如果之后另一条流数据到达且能 Join 到,则输出 +[L, R]
    • 随着 Watermark 推进,过期的数据会被删除并输出:
      • 左流过期输出 +[L, null]
      • 右流过期输出 +[null, R]

    3.3 示例

    — 直播间曝光日志表
    CREATE TABLE show_log_table (
    log_id BIGINT,
    show_params STRING,
    row_time AS CAST(CURRENT_TIMESTAMP AS TIMESTAMP(3)),
    WATERMARK FOR row_time AS row_time
    ) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '1',
    'fields.show_params.length' = '3',
    'fields.log_id.min' = '1',
    'fields.log_id.max' = '10'
    );

    — 直播间点击日志表
    CREATE TABLE click_log_table (
    log_id BIGINT,
    click_params STRING,
    row_time AS CAST(CURRENT_TIMESTAMP AS TIMESTAMP(3)),
    WATERMARK FOR row_time AS row_time
    ) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '1',
    'fields.click_params.length' = '3',
    'fields.log_id.min' = '1',
    'fields.log_id.max' = '10'
    );

    CREATE TABLE sink_table (
    s_id BIGINT,
    s_params STRING,
    c_id BIGINT,
    c_params STRING
    ) WITH (
    'connector' = 'print'
    );

    — 曝光流 Inner Join 点击流(曝光后 4 小时内的点击)
    INSERT INTO sink_table
    SELECT
    show_log_table.log_id AS s_id,
    show_log_table.show_params AS s_params,
    click_log_table.log_id AS c_id,
    click_log_table.click_params AS c_params
    FROM show_log_table
    INNER JOIN click_log_table
    ON show_log_table.log_id = click_log_table.log_id
    AND show_log_table.row_time BETWEEN
    click_log_table.row_time INTERVAL '4' HOUR
    AND click_log_table.row_time;

    3.4 关键要点

    • Interval Join 仅支持 append-only 表,即只追加不更新的数据流
    • 两表都需要有可比较的时间属性(事件时间或处理时间),且类型必须一致
    • 时间区间条件限定了一条流中的数据只会与另一条流中时间范围匹配的数据进行关联,这有效限制了 State 的存储范围
    • 相对于 Regular Join,Interval Join 不会产生回撤流(因为数据过期时发送的是 insert 而非 retract)
    • Flink 1.20 引入了 Early Fire 支持,可通过 Hint 指定提前触发的策略

    3.5 Interval Join vs Regular Join 选型对比

    对比维度Regular JoinInterval Join
    适用表类型 任意类型(append/upsert/retract) 仅 append-only 表
    时间属性要求 不强制 必须有且类型一致
    State 增长 无限增长,需配 TTL 自动过期,受水印控制
    回撤流 外连接会产生 不产生
    结果完整性 关联全部历史数据 仅关联时间窗口内的数据

    四、Event-Time Temporal Join(事件时间时态关联)

    4.1 什么是 Temporal Join

    Temporal Join(时态表关联)是一种基于时间属性的关联操作,它允许一条流记录根据其时间戳,去查询另一张“随时间变化的表”在该时刻的状态。简而言之:Temporal Join 回答的是“当时是什么”而非“现在是什么”的问题。

    在流处理中,常遇到两类需求:

    场景问题解决方案
    流 + 动态维表 维度数据会更新(如用户地址变更) 关联当前最新维表(Lookup Join)
    流 + 变化历史表 需关联事件发生时的维度状态(如订单创建时的汇率) 关联历史版本维表(Event-Time Temporal Join)

    因此,Temporal Join 包含两种类型:

    • Event-Time Temporal Join(事件时间时态关联):查询“过去某个时刻”的维表状态 → 本章重点
    • Processing-Time Temporal Join(处理时间时态关联):查询“当前最新”的维表状态 → 详见第五章 Lookup Join

    为了理解 Event-Time Temporal Join,我们首先需要理解一个关键概念:Versioned Table(版本表)。

    4.2 Versioned Table(版本表)

    4.2.1 概念简介

    Versioned Table(版本表)是一张随时间演化的表,其中每条主键记录可能有多个版本,每个版本有一个有效时间区间 [valid_from, valid_to),用于表示该版本在哪个时间段内生效。

    例如,汇率变化流在 Flink 内部会被物化为:

    currencyratevalid_fromvalid_to
    USD 7.20 2024-06-01 10:00:00 2024-06-01 11:00:00
    USD 7.25 2024-06-01 11:00:00 +∞

    关键特性:

    • 同一主键(currency)多版本共存
    • 时间区间为左闭右开:[valid_from, valid_to)
    • 最新版本的 valid_to = +∞
    4.2.2 数据模型:从 Changelog 到 Versioned Table

    Versioned Table 并非直接创建,而是由 Changelog Stream(变更日志流)构建而来,来源包括 Debezium、Canal、Flink CDC 等。

    物理输入——Changelog 流:

    opproduct_idpriceupdate_time
    +I P001 100 2024-06-01 09:00
    +U P001 120 2024-06-01 11:00
    +U P001 90 2024-06-01 15:00

    逻辑视图——Flink 内部构建的 Versioned Table:

    product_idpricevalid_fromvalid_to
    P001 100 2024-06-01 09:00:00 2024-06-01 11:00:00
    P001 120 2024-06-01 11:00:00 2024-06-01 15:00:00
    P001 90 2024-06-01 15:00:00 +∞
    4.2.3 如何构建 Versioned Table

    Flink 不提供显式 DDL 创建 Versioned Table,而是通过以下方式隐式构建。只要同时满足以下三个条件,Flink 就会自动将该表视为 Versioned Table:

  • 输入必须是 Changelog 流(有完整的 INSERT/UPDATE/DELETE 操作类型)
  • 定义主键(PRIMARY KEY)
  • 定义事件时间属性(Event Time)和 WATERMARK
  • — 定义 Changelog 流(即 Versioned Table 的来源)
    CREATE TABLE product_prices (
    product_id STRING,
    price DECIMAL(10, 2),
    update_time TIMESTAMP(3),
    — 必须定义事件时间 + watermark
    WATERMARK FOR update_time AS update_time INTERVAL '5' SECOND,
    — 必须声明主键
    PRIMARY KEY (product_id) NOT ENFORCED
    ) WITH (
    'connector' = 'kafka',
    'format' = 'debezium-json',
    'topic' = 'mysql.products.price_updates'
    );

    ⚠️ 重要提示:Flink 会在 State Backend(如 RocksDB)中维护 Versioned Table 的全量历史。每个主键的所有版本按时间排序存储,状态大小 = 主键数量 × 平均版本数,可能非常大!必须配置 State TTL 防止 OOM。

    4.2.4 Hive 拉链表与 Versioned Table 的区别

    Flink 无法直接将 Hive 拉链表识别为 Versioned Table,原因如下:

    • Hive 拉链表是静态快照表,不是流
    • 无操作类型(op 字段)
    • 无事件时间字段(均为业务时间,非变更时间)
    • Flink 读 Hive 表默认是 Batch 模式,无法产生 Changelog

    如需将 Hive 拉链表转变为 Versioned Table,推荐的做法是:通过 Flink CDC 直接同步源系统,不依赖 Hive 拉链表,而是直接监听源数据库(如 MySQL)的 Binlog,用 Flink CDC 将用户表变更实时捕获为 Changelog 流。

    4.3 Event-Time Temporal Join 实战

    4.3.1 用途与语法

    Event-Time Temporal Join 用于关联事件发生时的历史维度状态,精确回溯历史状态,是高级实时数仓的核心能力。其核心语法为:

    JOIN versioned_table FOR SYSTEM_TIME AS OF <event_time_attribute>

    语法源自 SQL:2011 标准中定义的时态表关联语法。

    4.3.2 完整示例

    场景:订单流关联汇率变化流,查询下单时点的汇率来精确计算金额。

    — Step 1: 定义汇率变化 Changelog 流(Versioned Table)
    CREATE TABLE currency_rates (
    currency STRING,
    rate DECIMAL(10, 4),
    update_time TIMESTAMP(3),
    WATERMARK FOR update_time AS update_time INTERVAL '5' SECOND,
    PRIMARY KEY (currency) NOT ENFORCED
    ) WITH (
    'connector' = 'kafka',
    'value.format' = 'debezium-json',
    'topic' = 'mysql.finance.currency_rates'
    );

    — Step 2: 定义订单流
    CREATE TABLE order_stream (
    order_id BIGINT,
    currency STRING,
    amount DECIMAL(10, 2),
    row_time TIMESTAMP(3),
    WATERMARK FOR row_time AS row_time INTERVAL '5' SECOND
    ) WITH (
    'connector' = 'kafka',
    'value.format' = 'json',
    'topic' = 'orders'
    );

    — Step 3: Event-Time Temporal Join
    SELECT
    o.order_id,
    o.amount,
    r.rate,
    o.amount * r.rate AS amount_cny
    FROM order_stream AS o
    JOIN currency_rates FOR SYSTEM_TIME AS OF o.row_time AS r
    ON o.currency = r.currency;

    4.3.3 关联过程详解

    Flink 内部将 Changelog 流物化为按时间版本化的表:

    currencyratevalid_fromvalid_to
    USD 7.20 10:00:00 11:00:00
    USD 7.25 11:00:00 +∞

    假设一条订单到达,order_time = 2024-06-01 10:30:00:

  • 提取 order_time = 10:30:00
  • 从 Versioned Table 中查找满足条件的版本:WHERE currency = 'USD'
    AND valid_from <= '10:30:00'
    AND '10:30:00' < valid_to

  • 匹配到 rate = 7.20(因为 valid_from ≤ 10:30 < valid_to)
  • 输出结果
  • 4.3.4 重要约束
    • 事件时间 Temporal Join 由左右两侧的 Watermark 触发;请确保 Join 的两边都正确设置了 Watermark
    • 事件时间 Temporal Join 要求主键包含在关联条件中
    • 维表必须是 Changelog 流,需来自 CDC(如 Debezium + Kafka),数据包含完整的变更日志(INSERT/UPDATE/DELETE)

    五、Lookup Join(维表 Join)

    5.1 什么是 Lookup Join

    Lookup Join 指的是:从一条流数据(主表 / 事实表)出发,根据某个关联字段(如 user_id),实时去外部系统(维表)查询对应的维度信息(如 user_name、city),以丰富流数据的过程。

    这本质上是一种**“流 + 维表”的实时关联**(Real-time Enrichment)。每条流记录触发一次对外部系统的点查(Lookup)。

    Lookup Join 通常用于用从外部系统查询的数据来丰富表。Join 要求一个表具有处理时间属性,另一个表由支持 Lookup 的源连接器提供支持。两个表之间还需要一个强制的等值连接条件。

    5.2 哪些系统可以作为维表

    并非所有数据源都能作为维表用于实时 Lookup Join,只有那些支持**“按主键快速点查”**的外部系统才被称为“可 Lookup 的”。

    可 Lookup 的系统不可 Lookup 的系统
    MySQL / PostgreSQL(JDBC) Hive(存储在 HDFS,无随机读)
    HBase Kafka Topic(日志流,无法按 Key 随机读取)
    Redis 普通文件(CSV / Parquet,需全表扫描)
    Elasticsearch 未建索引的数据库大表
    MongoDB

    判断标准:能否在 < 100ms 内根据主键返回单条记录?如果不能,就不适合作为维表进行 Lookup Join。

    5.3 Processing-Time Temporal Join

    根据 Flink 官方文档,基于处理时间的 Temporal Join(Processing-Time Temporal Join)等价于 Lookup Join。

    语法上,两者都使用 FOR SYSTEM_TIME AS OF:

    — 处理时间 Temporal Join / Lookup Join
    JOIN dim_table FOR SYSTEM_TIME AS OF <processing_time_attribute>

    Flink 内部对这一语法有两种处理方式:

    处理方式适用场景
    时态表函数(Temporal Table Function) 维表是 Changelog 流(CDC)
    Lookup Source 维表是外部系统(JDBC/HBase/Redis 等)

    本文聚焦后者,即维表为外部系统的情形。

    5.4 Lookup Join 的技术原理

    Flink 通过专门的 Lookup Join 算子实现维表关联,其内部机制主要包括:

    5.4.1 同步 Lookup(不推荐)

    每条流记录到来时,阻塞等待维表查询返回后再继续处理下一条。这种方式吞吐低,容易产生背压。

    5.4.2 异步 Lookup(Async Lookup,推荐)

    使用 Flink Async I/O API 实现,多条记录可以并发查询,结果通过回调处理。这种方式具有高吞吐、低延迟的优势。

    维表 DDL 需开启异步查询:

    'lookup.async' = 'true'

    5.4.3 缓存优化

    为减少对外部系统的重复查询,Lookup Join 支持 LRU 缓存机制:

    'lookup.cache.max-rows' = '300000',
    'lookup.cache.ttl' = '20 min'

    5.4.4 Hash Lookup Join(Flink 1.18+)

    为了进一步提升 Lookup Join 的缓存命中率,Flink 引入了 Hash Lookup Join 机制。通过将相同 Lookup Key 的数据路由到同一个 Task 实例,可以显著提高缓存命中率:

    SELECT /*+ SHUFFLE_HASH('Customers') */
    o.order_id, o.total, c.country, c.zip
    FROM Orders AS o
    JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c
    ON o.customer_id = c.id;

    5.5 Flink 对 PROCTIME() 的隐式处理

    即使你在 SQL 中没有显式写出 PROCTIME(),Flink 内部也会自动完成处理:

    — 用户写的 SQL(看似没有 proc_time)
    SELECT o.*, u.name
    FROM orders AS o
    JOIN mysql_users AS u
    ON o.user_id = u.user_id;

    Flink 内部执行以下步骤:

  • 检测到 mysql_users 是可 Lookup 的维表(通过 Connector 类型判断)
  • 自动将此 Join 重写为 Temporal Join
  • 隐式注入一个处理时间属性
  • 等价于:
  • SELECT o.*, u.name
    FROM orders AS o
    JOIN mysql_users FOR SYSTEM_TIME AS OF PROCTIME() AS u
    ON o.user_id = u.user_id;

    5.6 需要显式写出 PROCTIME() 的场景

    场景一:多次关联不同维表

    若不显式定义,Flink 会为每个 Join 生成独立的 PROCTIME(),可能导致两次查询的时间不一致:

    — 显式定义 proc_time,避免歧义
    CREATE TABLE orders (
    order_id STRING,
    user_id STRING,
    product_id STRING,
    ptime AS PROCTIME() — 显式声明
    );

    SELECT
    o.*,
    u.name,
    p.category
    FROM orders AS o
    JOIN users FOR SYSTEM_TIME AS OF o.ptime AS u ON o.user_id = u.user_id
    JOIN products FOR SYSTEM_TIME AS OF o.ptime AS p ON o.product_id = p.id;

    场景二:与 Event-Time 流混合使用

    当作业中同时使用了 Event-Time Temporal Join 和 Lookup Join 时,需要显式给出处理时间字段,否则 Flink 不清楚用哪个时间做 Lookup。

    场景三:代码可读性与调试

    显式写出 FOR SYSTEM_TIME AS OF o.proc_time 能让代码意图更清晰,便于团队协作和维护。

    5.7 PROCTIME() 的工作机制

    PROCTIME() 在 Lookup Join 中扮演着时间锚点(Time Anchor)的角色。其核心作用如下:

  • 为每条流记录打上处理时间戳:作为“查询发起时刻”的标记
  • 提供 Temporal Join 的语法合法性:Flink 要求所有 Temporal Join 必须指定一个时间属性
  • 确定“查询维表的时刻”:维表的最新值是“记录被处理那一刻”的值
  • 注意:在 Lookup Join 中,proc_time 的具体数值其实不重要。重要的是它作为一个“触发信号”,告诉 Flink:“现在该去查维表了”。因为维表只返回当前最新值,不关心历史版本。

    5.8 如何判断一个 Connector 是否支持 Lookup

    查看 Flink 官方文档中该 Connector 是否支持以下特性:

  • 支持 PRIMARY KEY 声明
  • 支持 lookup 相关配置参数,如:
    • 'lookup.async'
    • 'lookup.cache.max-rows'
    • 'lookup.cache.ttl'
    • 'lookup.batch-size'
  • 文档明确说明可用于 “Temporal Table Join” 或 “Lookup Join”
  • 例如,Flink JDBC Connector 文档明确写道:“The JDBC table can be used for lookup joins with temporal tables.”。而 Hive Connector 文档则从未提及 Lookup Join 支持。

    5.9 为什么 Hive 不能做维表

    这是一个常见的误区。Hive 表不能作为维表进行 Lookup Join 的核心原因:

  • 存储位置:Hive 表数据存储在 HDFS,无随机读能力
  • 查询引擎:需启动 MapReduce / Tez / Spark 作业,延迟高达秒至分钟级,无法做到毫秒级响应
  • Flink 设计:Flink 的 HiveCatalog 仅用于元数据读取和流式写入,不支持 Lookup Join
  • 若需使用 Hive 中的维度数据,应先将其同步到可 Lookup 的系统中:

    • Hive → HBase(通过 Sqoop / Spark)
    • Hive → MySQL(通过 DataX)
    • Hive → Redis(通过自定义程序)

    5.10 完整示例

    以下示例使用曝光用户日志流关联 Redis 用户画像维表:

    — 曝光用户日志流(事实表)
    CREATE TABLE show_log (
    log_id BIGINT,
    `timestamp` AS CAST(CURRENT_TIMESTAMP AS TIMESTAMP(3)),
    user_id STRING,
    proctime AS PROCTIME()
    ) WITH (
    'connector' = 'datagen',
    'rows-per-second' = '10',
    'fields.user_id.length' = '1',
    'fields.log_id.min' = '1',
    'fields.log_id.max' = '10'
    );

    — Redis 用户画像维表
    CREATE TABLE user_profile (
    user_id STRING,
    age STRING,
    sex STRING
    ) WITH (
    'connector' = 'redis',
    'hostname' = '127.0.0.1',
    'port' = '6379',
    'format' = 'json',
    'lookup.cache.max-rows' = '500',
    'lookup.cache.ttl' = '3600',
    'lookup.max-retries' = '1'
    );

    — 结果表
    CREATE TABLE sink_table (
    log_id BIGINT,
    `timestamp` TIMESTAMP(3),
    user_id STRING,
    proctime TIMESTAMP(3),
    age STRING,
    sex STRING
    ) WITH (
    'connector' = 'print'
    );

    — Lookup Join 查询
    INSERT INTO sink_table
    SELECT
    s.log_id,
    s.`timestamp`,
    s.user_id,
    s.proctime,
    u.sex,
    u.age
    FROM show_log AS s
    LEFT JOIN user_profile FOR SYSTEM_TIME AS OF s.proctime AS u
    ON s.user_id = u.user_id;

    5.11 维表 Join 注意事项

    以下是生产环境中使用维表 Join 需要特别注意的关键点:

  • 维表必须支持高效点查

    • 禁止使用 Hive、Kafka、普通文件作为维表
    • 必须使用 OLTP 或 KV 存储:MySQL、HBase、Redis、Elasticsearch 等
    • 验证标准:单次查询 P99 延迟 < 50ms
  • 主键必须匹配

    • 维表 DDL 中必须声明 PRIMARY KEY
    • Join 条件必须使用该主键字段
    • 否则 Flink 无法生成高效的查询计划
  • Checkpoint 超时必须 > Lookup 超时

    • Flink 要求在一次 Checkpoint 触发后,所有正在处理的异步请求必须在 Checkpoint 超时前完成
    • 如果 Lookup 的超时时间比 Checkpoint 超时还长,Checkpoint 就会因“等待异步操作完成”而超时失败
    • 推荐:Checkpoint 超时时间至少是 Lookup 超时时间的 2-4 倍
  • 维表延迟问题:对于会实时新建/更新的维表,应建立数据延迟监控机制,防止流表数据先于维表数据到达,导致关联不到维表数据

  • 结果时效性问题:同一条流数据关联的维表结果不会随维表更新而更新。如果维表数据发生变化,已关联到的结果数据不会再同步更新。在评估实时任务的准确性时需要考虑这一点

  • 5.12 维表 Join 优化

    作业配置优化

    — Checkpoint 必须 > lookup timeout
    SET 'execution.checkpointing.interval' = '2 min';
    SET 'execution.checkpointing.timeout' = '5 min';

    — 异步缓冲区防背压
    SET 'table.exec.async-lookup.buffer-capacity' = '500';

    维表连接器配置优化

    — 启用异步,设置超时和重试时长
    'lookup.async' = 'true',
    'lookup.async.timeout' = '20 s',
    'lookup.max-retries' = '2',

    — 启用缓存,减少重复查询
    'lookup.cache.max-rows' = '300000',
    'lookup.cache.ttl' = '20 min',

    — 开启批量查询,减少 RPC 调用次数
    'lookup.batch-size' = '400',
    'lookup.buffer-timeout' = '150 ms'

    优化维度
    优化手段效果适用场景
    异步 Lookup 提升吞吐,降低延迟 所有场景(强烈推荐)
    缓存(LRU) 减少外部查询次数 维表数据变化不频繁的场景
    Hash Lookup Join 提高缓存命中率 高并发场景
    增加 JVM 内存 + 缓存记录数 提升缓存容量 内存资源充足的场景
    维表建立索引 加快单次查询速度 数据库维表

    六、State TTL 管理与配置

    6.1 为什么需要 State TTL

    在 Flink 流处理中,State(状态)是保证计算正确性的关键。Regular Join 需要将 Join 两侧的数据永久保存在状态中,资源消耗可能无限增长。如果不加以控制,State 膨胀最终会导致:

    • TaskManager 内存溢出(OOM)
    • Checkpoint 耗时过长甚至超时失败
    • 作业恢复时间过长

    State TTL(Time-To-Live,生存时间)是控制状态大小的核心手段。

    6.2 State TTL 的配置层级

    作业级别 TTL(Flink 1.17 及之前)

    在 Flink 1.17 及更早版本中,State TTL 只能以作业粒度进行配置:

    SET 'table.exec.state.ttl' = '86400000'; — 1 天(毫秒)

    这意味着作业中所有有状态的算子共用同一个 TTL 值。但实际场景中,不同 Join 对状态保留时间的要求差异很大——例如订单表的状态可能需要保留 30 天,而用户点击日志可能只需保留 1 天。

    算子级别 TTL(Flink 1.18+)

    自 Flink 1.18 起,通过 FLIP-292 引入了算子级别的 State TTL 配置支持,允许通过 Compiled JSON Plan 为不同算子设置不同的 TTL 值。

    此外,FLIP-373 进一步提出了通过 SQL Hint 方式为 Regular Join 和 Group Aggregate 配置不同 TTL 的方案:

    — 使用 SQL Hint 为不同表设置不同 TTL
    SELECT /*+ STATE_TTL('orders' = '1d', 'customers' = '20d') */
    *
    FROM orders
    LEFT OUTER JOIN customers
    ON orders.o_custkey = customers.c_custkey;

    — 也可以使用表别名
    SELECT /*+ STATE_TTL('o' = '12h', 'c' = '3d') */
    *
    FROM orders o
    LEFT OUTER JOIN customers c
    ON o.o_custkey = c.c_custkey;

    6.3 State TTL 配置建议

    Join 类型建议 TTL原因
    Regular Join 根据业务需要设置(如 1d-30d) 需权衡正确性与资源
    Interval Join 可不设,或设较大值 Watermark 自动清理
    Event-Time Temporal Join 根据维表变更频率设置 需保留维表历史版本
    Lookup Join 无需 State TTL(维表不在 State 中) 外部查询

    ⚠️ 警告:为 Regular Join 设置过短的 TTL 会导致早期数据被清除,可能影响 Join 结果的正确性。需要在状态大小与结果完整性之间做出权衡。

    七、列转行(Array Expansion & Table Function)

    7.1 Array Expansion(数组展开)

    Array Expansion 是 Flink SQL 中的一种列转行操作,类似于 Hive 中的 explode 函数,用于将单行中的数组类型字段展开为多行。

    — 示例:将标签数组展开为多行
    SELECT
    user_id,
    tag
    FROM users
    CROSS JOIN UNNEST(tags) AS t(tag);

    7.2 Table Function(表函数)

    Table Function 是用户自定义函数(UDTF),将一行输入映射为多行输出,支持 Inner Join 和 Left Outer Join:

    — 使用表函数进行列转行
    SELECT
    order_id,
    item
    FROM orders
    LEFT JOIN LATERAL TABLE(split_to_items(order_list)) AS t(item) ON TRUE;

    Lateral Join 将表与表值函数的结果连接,左表的每一行都与表值函数相应调用产生的所有行相连接,要求 ON 子句中有固定的 TRUE 连接条件。

    八、Join 类型全景对比与选型指南

    8.1 Join 类型对比总表

    维度Regular JoinInterval JoinEvent-Time Temporal JoinLookup Join
    关联对象 流 ↔ 流 流 ↔ 流 流 ↔ 维表(Changelog) 流 ↔ 外部维表
    时间语义 不限 事件时间 / 处理时间 事件时间 处理时间
    State 增长 无限(需 TTL) 自动过期(Watermark) 需 TTL(保留维表历史) 维表不在 State 中
    回撤流 外连接有
    语法 JOIN … ON JOIN … ON … AND time JOIN … FOR SYSTEM_TIME AS OF JOIN … FOR SYSTEM_TIME AS OF
    维表来源 任意流 任意流 CDC(Kafka + Debezium) MySQL/HBase/Redis 等
    典型场景 实时数仓宽表 Join 曝光-点击关联分析 汇率/价格的历史回溯 用户画像实时补全

    8.2 选型决策流程

    是否需要关联事件发生时的历史状态?
    ├── 是 → Event-Time Temporal Join
    │ └── 前提:维表是 CDC Changelog 流
    └── 否 → 维表是外部系统还是流?
    ├── 外部系统(MySQL/Redis/HBase)
    │ └── Lookup Join(异步 + 缓存)
    └── 另一条流
    ├── 需要避免回撤流,且有时间窗口约束?
    │ └── Interval Join
    └── 需要完整的 Join 语义?
    └── Regular Join(别忘了配 State TTL)

    九、生产环境最佳实践

    9.1 Join 作业优化清单

  • Join 顺序:在 FROM 子句中,将更新频率最低的表放在前面,频率最高的放在后面
  • State TTL:为每个 Join 算子设置合理的 TTL,避免状态无限增长
  • 异步 Lookup:维表 Join 必须开启 lookup.async = true
  • 缓存配置:根据维表大小设置合理的 lookup.cache.max-rows 和 lookup.cache.ttl
  • Checkpoint 超时:确保 Checkpoint 超时时间 > Lookup 超时时间的 2-4 倍
  • Hash Lookup Join:高并发场景考虑开启 SHUFFLE_HASH
  • 资源规划:Regular Join 和 Temporal Join 需要充足的堆外内存(RocksDB)
  • 监控告警:监控 State 大小、Checkpoint 耗时、维表查询延迟
  • 9.2 常见踩坑点

    踩坑点表现解决方案
    Regular Join 不设 TTL State 无限增长,最终 OOM 设置 table.exec.state.ttl
    TTL 设置过短 Join 结果不完整 评估业务窗口,设置合理 TTL
    Lookup 超时 > Checkpoint 超时 Checkpoint 频繁超时失败 调整超时比例或优化维表性能
    用 Hive 做维表 查询延迟极高,作业卡住 同步到 HBase/Redis/MySQL
    Interval Join 两表时间类型不一致 语法错误 统一使用事件时间或处理时间
    未开启异步 Lookup 吞吐量极低 设置 lookup.async = true

    十、总结

    Flink SQL 提供了丰富的 Join 能力,涵盖了流与流关联、流与维表关联、列转行等多种场景。选择合适的 Join 类型是实时数仓建设中的关键决策:

    • Regular Join 最为灵活通用,但需要严格控制 State 大小
    • Interval Join 通过时间窗口约束避免了回撤流,适合有明确时间范围的关联场景
    • Event-Time Temporal Join 是精确历史回溯的核心能力,需要维表以 CDC Changelog 流形式提供
    • Lookup Join 是最常用的流与外部维表关联方式,开启异步和缓存是基本配置

    在实际生产中,理解每种 Join 的内部机制和适用边界,合理配置 State TTL 和维表参数,才能在保证正确性的同时实现高性能的实时数据处理。

    赞(0)
    未经允许不得转载:171主机测评 » Flink 系列第23篇:Flink SQL 多流 Join 指南 —— 从原理到调优,一文吃透
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址