【发布时间】:2021-02-27 06:32:43
【问题描述】:
我正在尝试使用 Spark 2.3.0 实现流-流连接玩具
条件匹配时流连接工作正常,但即使使用leftOuterJoin,条件不匹配时也会丢失左流值。
提前致谢
这是我的源代码和数据,基本上,我正在创建两个套接字,一个是 9999 作为右流源,一个是 9998 作为左流源。
val spark = SparkSession
.builder
.appName("StreamStream")
.master("local")
.getOrCreate()
import spark.implicits._
spark.sparkContext.setLogLevel("ERROR")
val s9999: DataFrame = spark
.readStream
.format("socket")
.option("host", "localhost")
.option("port", 9999)
.load()
val s9999Dataset: Dataset[S9999] = s9999
.map(line => {
val strings = line.get(0).toString.split(",")
val id = strings(0).toInt
val time = Timestamp.valueOf(strings(1))
S9999(id, time)
})
.withWatermark("timestamp99", "30 seconds")
val s9998Dataset: Dataset[S9998] = spark
.readStream
.format("socket")
.option("host", "localhost")
.option("port", 9998)
.load()
.map(line => {
val strings = line.get(0).toString.split(",")
val id = strings(0).toInt
val time = Timestamp.valueOf(strings(1))
S9998(id, time)
})
val resultDataset = s9998Dataset
.join(s9999Dataset,
joinExprs = expr(
"""
id99 = id98 AND
timestamp98 >= timestamp99 AND
timestamp98 <= timestamp99 + interval 6 seconds
"""),
joinType = "leftOuter")
val streamingQuery = resultDataset
.writeStream
.outputMode("append")
.format("console")
.start()
streamingQuery.awaitTermination()
}
case class S9999(id99: Int, timestamp99: Timestamp)
case class S9998(id98: Int, timestamp98: Timestamp)
样本数据:
左插座:
1,2011-10-02 18:50:20.123
2,2011-10-02 18:50:25.123
3,2011-10-02 18:50:30.123
4,2011-10-02 18:50:35.123
5,2011-10-02 18:50:40.123
6,2011-10-02 18:50:45.123
7,2011-10-02 18:50:50.123
8,2011-10-02 18:50:55.123
9,2011-10-02 18:51:00.123
10,2011-10-02 18:51:05.123
11,2011-10-02 18:51:10.123
12,2011-10-02 18:51:15.123
13,2011-10-02 18:51:20.123
14,2011-10-02 18:51:25.123
15,2011-10-02 18:51:30.123
正确的流数据:
1,2011-10-02 18:50:20.123
3,2011-10-02 18:50:30.123
7,2011-10-02 18:50:50.123
8,2011-10-02 18:50:55.123
9,2011-10-02 18:51:00.123
13,2011-10-02 18:51:20.123
14,2011-10-02 18:51:25.123
15,2011-10-02 18:51:30.123
【问题讨论】:
-
@philipxy 查看了我的正确流的示例数据,水印设置为 30 秒,但我在这里发布的数据大约为 1 分钟,如果我放一些稍后的时间戳,它也不起作用,例如2011-10-02 18:53:30.123。 “简而言之,如果正在加入的两个输入流中的任何一个在一段时间内没有接收到数据,那么外部(两种情况,左或右)输出可能会延迟。”正如 Bah91 建议的那样,但是什么时候会发出空值?
-
LEFT JOIN ON 返回 INNER JOIN ON 行 UNION ALL 由 NULL 扩展的不匹配的左表行。直到获得所有正确的表行之后,才能知道给定的左表行是否不匹配。在已知某个流完成之前,不能发出空扩展行。你不应该给99加水印吗?你的循环是由左/99 行驱动的,所以你不应该在所有右/98 行之后加入吗?我没有看到您这样做,但我所知道的只是流式传输、加入、此代码和左加入手册部分。 (从您的帖子中不清楚您是否也有 INNER JOIN 的问题。)
-
@philipxy水印右边是必填的,98是左边流,99是右边;内部连接很好。
标签: apache-spark left-join spark-structured-streaming apache-spark-2.3