【发布时间】:2019-12-19 21:02:38
【问题描述】:
我正在学习/体验 Flink,我观察到 DataStream 连接的一些意外行为,并想了解正在发生的事情......
假设我有两个流,每个流有 10 条记录,我想加入 id 字段。让我们假设一个流中的每条记录在另一个流中都有一个匹配的记录,并且每个流中的 ID 都是唯一的。假设我必须使用全局窗口(要求)。
使用 DataStream API 加入(我在 Scala 中的简化代码):
val stream1 = ... // from a Kafka topic on my local machine (I tried with and without .keyBy)
val stream2 = ...
stream1
.join(stream2)
.where(_.id).equalTo(_.id)
.window(GlobalWindows.create()) // assume this is a requirement
.trigger(CountTrigger.of(1))
.apply {
(row1, row2) => // ...
}
.print()
结果:
- 一切都按预期打印,第一个流中的每条记录都与第二个流中的一条记录相结合。
但是:
- 如果我将其中一条记录(例如,带有更新的字段)从一个流重新发送到该流,则会发出两个重复的连接事件????
- 如果我重复该操作(有或没有更新的字段),我将得到 3 个已发出的事件,然后是 4、5 等等...... ????
Flink 社区中有人能解释一下为什么会这样吗?我本来预计每次只会发出 1 个事件。是否可以通过全局窗口来实现这一点?
相比之下,Flink Table API 在同一场景中的表现与预期一致,但对于我的项目,我对 DataStream API 更感兴趣。
Table API 示例,按预期工作:
tableEnv
.sqlQuery(
"""
|SELECT *
| FROM stream1
| JOIN stream2
| ON stream1.id = stream2.id
""".stripMargin)
.toRetractStream[Row]
.filter(_._1) // just keep the inserts
.map(...)
.print() // works as expected, after re-sending updated records
谢谢,
尼古拉斯
【问题讨论】:
-
您能否更具体地说明“如果我将其中一条记录(例如,带有更新的字段)从一个流重新发送到该流,则会发出两个重复的连接事件。 "能否举个简单的例子。
-
例如,如果流 1 只有一个记录 KeyValueRecord(1, 10),流 2 只有一个 KeyValueRecord(1, 42),我的应用程序将打印 [KeyValueRecord(1, 10), KeyValueRecord (1, 42)],因为两条记录具有相同的键“1”。如果稍后,我将新记录 KeyValueRecord(1, 11) 推送到流 1,我的应用程序不仅会再次打印之前的 [KeyValueRecord(1, 10), KeyValueRecord(1, 42)],还会打印 [KeyValueRecord(1, 11), KeyValueRecord(1, 42)] (我只期望后者)。如果我再次推送相同的记录,它会打印相同的内容,再加上 [KeyValueRecord(1, 11), KeyValueRecord(1, 42)] 再次等等...
标签: apache-flink