【问题标题】:Apache Flink: How to enable "upsert mode" for dynamic tables?Apache Flink:如何为动态表启用“upsert 模式”?
【发布时间】:2018-07-11 07:41:43
【问题描述】:

我已经在 Flink 文档和官方 Flink 博客中看到过多次提到基于唯一键的动态表的“upsert 模式”。但是,我没有看到任何关于如何在动态表上启用此模式的示例/文档。

例子:

  • Blog post:

    当通过更新模式在流上定义动态表时,我们可以在表上指定一个唯一键属性。在这种情况下,对键属性执行更新和删除操作。 更新模式如下图所示。

  • Documentation:

    转换为 upsert 流 的动态表需要一个(可能是复合的)唯一键

所以我的问题是:

  • 如何在 Flink 中为动态表指定唯一键属性?
  • 如何将动态表置于更新/更新插入/“替换”模式而不是追加模式?

【问题讨论】:

    标签: apache-flink flink-streaming flink-sql


    【解决方案1】:

    链接的资源描述了两种不同的场景。

    • blog post 讨论了 upsert DataStream -> Table 转换。
    • documentation 描述了反向 upsert Table -> DataStream 转换。

    以下讨论基于 Flink 1.4.0(2018 年 1 月)。

    Upsert DataStream -> Table 转化

    本机不支持通过键上的 upsert 将 DataStream 转换为 Table,但在路线图上。同时,您可以使用附加 Table 和具有用户定义聚合函数的查询来模拟此行为。

    如果您有一个追加 Table Logins 与架构 (user, loginTime, ip) 跟踪用户的登录,您可以将其转换为一个 upsert Table 键入 user 使用以下查询:

    SELECT user, LAST_VAL(loginTime), LAST_VAL(ip) FROM Logins GROUP BY user
    

    LAST_VAL 聚合函数是一个user-defined aggregation function,它始终返回最新的附加值。

    对 upsert DataStream -> Table 转换的本机支持基本上以相同的方式工作,但提供了更简洁的 API。

    Upsert Table -> DataStream 转化

    不支持将 Table 转换为 upsert DataStream。这也正确反映在文档中:

    请注意,将动态表转换为 DataStream 时,仅支持追加和收回流。

    我们故意选择不支持 upsert Table -> DataStream 转换,因为只有知道关键属性时才能处理 upsert DataStream。这些取决于查询,并不总是很容易识别。开发人员有责任确保正确解释关键属性。不这样做会导致错误的程序。为避免出现问题,我们决定不提供 upsert Table -> DataStream 转换。

    相反,用户可以将Table 转换为撤回DataStream。此外,我们支持 UpsertTableSink 将 upsert DataStream 写入外部系统,例如数据库或键值存储。

    【讨论】:

    • 非常有见地的答案,绝对为我解决了这个问题。谢谢@Fabian
    • 嗨@Fabian,LAST_VAL 是否已实施并在某处可用?
    • 不,它还不是内置函数。您需要将其实现为用户定义的聚合函数。但是,Flink 社区正在努力添加适当的 upsert 表支持。如果一切按计划进行,它应该包含在 1.6.0 版本中。
    • 嗨@FabianHueske - 感谢有关如何通过自定义 UDAF 将 upsert 流转换为动态表的提示。我在本地工作,当任何输入流发生变化时,3 路连接现在会触发撤回更新。你是否仍然计划为这种转换添加原生支持 - 我在 Flink 1.9 中看不到任何其他方式来做到这一点。谢谢安迪
    • 是的,此功能仍在路线图中,但现在社区对其他功能给予了更高的优先级。
    【解决方案2】:

    更新:从 Flink 1.9 开始,LAST_VALUEbuild-in aggregate functions 的一部分,如果我们使用 Blink planner(这是 Flink 1.11 以来的默认设置)。

    假设上面 Fabian Hueske 的回复中提到的 Logins 表存在,我们现在可以将其转换为 upsert 表,如下所示:

    SELECT 
      user, 
      LAST_VALUE(loginTime), 
      LAST_VALUE(ip) 
    FROM Logins 
    GROUP BY user
    

    【讨论】:

      【解决方案3】:

      Flink 1.8 仍然缺乏这样的支持。预计将来会添加这些功能:1) LAST_VAL 2) Upsert Stream Dynamic Table.

      ps。 LAST_VAL() 似乎不可能在 UDTF 中实现。聚合函数不提供附加的事件/过程时间上下文。阿里巴巴的 Blink 提供了 LAST_VAL 的替代实现,但它需要另一个字段来提供订单信息,而不是直接在 event/proc time 上。这使得 sql 代码很难看。 (https://help.aliyun.com/knowledge_detail/62791.html)

      我的 LAST_VAL(例如获取最新 ip)的变通解决方案类似于:

      1. concat(ts, ip) as ordered_ip
      2. MAX(ordered_ip) as ordered_ip
      3. extract(ordered_ip) as ip

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2019-07-25
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-06-12
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多