【问题标题】:Flink SQL Watermark Strategy After Join OperationJoin 操作后的 Flink SQL Watermark 策略
【发布时间】:2022-08-21 23:33:09
【问题描述】:

我的问题是我不能在 JOIN 操作之后使用 ORDER BY 子句。要重现问题,

CREATE TABLE stack (
    id INT PRIMARY KEY,
    ts TIMESTAMP(3),
    WATERMARK FOR ts AS ts - INTERVAL \'1\' SECONDS
) WITH (
  \'connector\' = \'datagen\',
  \'rows-per-second\' = \'5\',
  \'fields.id.kind\'=\'sequence\',
 \'fields.id.start\'=\'1\',
 \'fields.id.end\'=\'100\'
);

此表有水印策略和TIMESTAMP(3) *ROWTIME* 类型ts

Flink SQL> DESC stack;
+------+------------------------+-------+---------+--------+----------------------------+
| name |                   type |  null |     key | extras |                  watermark |
+------+------------------------+-------+---------+--------+----------------------------+
|   id |                    INT | FALSE | PRI(id) |        |                            |
|   ts | TIMESTAMP(3) *ROWTIME* |  TRUE |         |        | `ts` - INTERVAL \'1\' SECOND |
+------+------------------------+-------+---------+--------+----------------------------+
2 rows in set

但是,如果我将视图定义为简单的自联接

CREATE VIEW self_join AS (
SELECT l.ts, l.id, r.id
FROM stack as l INNER JOIN stack as r
ON l.id=r.id
);

它丢失了水印策略,但没有丢失类型,

Flink SQL> DESC self_join;
+------+------------------------+-------+-----+--------+-----------+
| name |                   type |  null | key | extras | watermark |
+------+------------------------+-------+-----+--------+-----------+
|   ts | TIMESTAMP(3) *ROWTIME* |  TRUE |     |        |           |
|   id |                    INT | FALSE |     |        |           |
|  id0 |                    INT | FALSE |     |        |           |
+------+------------------------+-------+-----+--------+-----------+
3 rows in set

我假设我们可以保留水印策略并在JOIN 操作之后使用ORDER BY,但事实并非如此。如何再次向VIEW 添加水印策略?

提前致谢。

    标签: apache-flink flink-sql


    【解决方案1】:

    每当 Flink SQL 在流模式下执行常规连接(没有任何时间约束的连接)时,结果不可能有水印。这反过来意味着您不能对结果进行排序或应用窗口化。

    为什么会这样,你能做些什么呢?

    背景

    Flink SQL 使用时间属性(在本例中为stack.ts)来优化状态保留。因为stack 流/表有一个时间属性,我们知道这个流将按时间或多或少地按顺序处理(元素被限制为最多1 秒乱序)。然后,这对必须保留多少状态才能执行排序此表之类的操作施加严格的限制——1 秒长的缓冲区就足够了。

    如果stack 没有定义时间属性(即定义了水印的时间戳字段),那么 Flink SQL 将拒绝对其进行排序(在流模式下),因为这样做需要保持一个无限量状态,并且不可能知道在发出第一个结果之前要等待多长时间。

    常规连接的结果不能有明确的水印策略

    任何类型的常规连接都要求 Flink 在其状态后端永久存储输入表的所有行(Flink 愿意尝试这样做)。但更重要的是,水印在结果上没有很好的定义,因为它可能是无序的没有限制。

    你可以做什么

    如果您将连接修改为interval jointemporal join,则结果仍然会有水印。例如,你可以这样做:

    CREATE VIEW self_join AS (
      SELECT l.ts, l.id, r.id
      FROM stack as l INNER JOIN stack as r
      ON l.id=r.id
      WHERE ls.ts BETWEEN r.ts - INTERVAL '1' MINUTE AND r.ts
    );
    

    或者你可以这样做:

    CREATE VIEW self_join AS (
      SELECT l.ts, l.id, r.id
      FROM stack as l INNER JOIN stack as r FOR SYSTEM_TIME AS OF r.ts
      ON l.id=r.id
    );
    

    在这两种情况下,Flink 的 SQL 引擎将能够保留比常规连接更少的状态,并且能够在输出流/表中产生水印。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-20
      • 1970-01-01
      相关资源
      最近更新 更多