【问题标题】:Flink: What are the best alternatives to using Python UDFs in MATCH_RECOGNIZE?Flink:在 MATCH_RECOGNIZE 中使用 Python UDF 的最佳替代方案是什么?
【发布时间】:2021-07-26 18:45:17
【问题描述】:

当我尝试在使用 Python UDF 的 SQL 查询中使用 MATCH_RECOGNIZE 时,我收到错误 Python Function can not be used in MATCH_RECOGNIZE for now.

例如,不支持以下内容:

SELECT T.aa as ta
FROM MyTable
MATCH_RECOGNIZE (
  ORDER BY proctime
  MEASURES
    A.a as aa,
    pyFunc(1,2) as bb
  PATTERN (A B)
  DEFINE
    A AS a = 1,
    B AS b = 'b'
) AS T

这引发了几个问题:

  1. 为什么 Blink 规划器需要支持 Python 函数?

  2. 在文档中哪里可以找到这种缺乏支持的情况?关于这个特性的docs 没有提到 Python。预计我会通过validation tests解析吗?

  3. (主要问题) MATCH_RECOGNIZE 的最佳替代方案是用户定义的表聚合 Python 函数吗?我只想按顺序查找两个事件(在一小时内)。我知道我可以通过自加入来做到这一点,但我想看看是否有更高效/更干净的可能性。

【问题讨论】:

    标签: python apache-flink match-recognize


    【解决方案1】:

    作为在 measure 子句中无法使用 Python UDF 的一种解决方法,您似乎可以从 MATCH_RECOGNIZE 生成作为 UDF 输入所需的数据,然后在后续步骤中应用 UDF。

    类似这样的:

    SELECT
      T.aa AS ta, 
      pyFunc(T.one, T.two) AS tb
    FROM MyTable
    MATCH_RECOGNIZE (
      ORDER BY proctime
      MEASURES
        A.a AS aa,
        1 AS one,
        2 AS two
      PATTERN (A B)
      DEFINE
        A AS a = 1,
        B AS b = 'b'
    ) AS T
    

    如果您决定改用这种方法,使用对时间属性有间隔约束的自联接应该会产生一个有效的计划。

    【讨论】:

    • 我想我会倾向于使用自连接,因为我不太清楚如何在 DEFINE 部分过滤没有 python UDF 的复杂序列。但很高兴这也是可取的,谢谢!
    • 抱歉,这只是切线相关,但我知道 LATERAL TABLE 用于 SQL 中的表函数,表聚合函数是否有 SQL 等效项?似乎文档只显示 flatAggregate。
    • 对于那些好奇的更新,Flink SQL 不支持表聚合函数。请改用普通的聚合函数。
    猜你喜欢
    • 2017-06-11
    • 1970-01-01
    • 2021-01-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-01-19
    • 1970-01-01
    • 2022-11-23
    相关资源
    最近更新 更多