链接的资源描述了两种不同的场景。
以下讨论基于 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 写入外部系统,例如数据库或键值存储。