【发布时间】:2021-02-17 14:21:20
【问题描述】:
我是使用 Scala 的 Apache Spark 的新手。我可以使用以下命令将表加入流:
Updated_DF = Inbound_DF.join(colToAdd, colToAdd("key") <=> Inbound_DF("key"), "left")
.withColumnRenamed("Data_DF","site").drop("Id","key")
现在我想检查colToAdd("key") 和Inbound_DF("key") 是否匹配并且加入是否成功。例如,colToAdd:
Id key Data_DF
S31 S3 {"name":"nick","region":"IN"}
S21 S2 {"name":"john","region":"CA"}
S11 S1 {"name":"ashley","region":"CA"}
S51 S5 {"name":"bella","region":"UK"}
S41 S4 {"name":"kumar","region":"In"}
S6 S6 {"name":"ben","region":"US"}
P11 P1 {"name":"MKD","region":"UAE"}
P21 P2 {"name":"ahmad","region":"UAE"}
来自传入流的消息如下所示:
cusId key item price
1897 S2 book 54
加入后,更新后的消息应如下所示:
cusId key item price site
1897 S2 book 54 {"name":"john","region":"CA"}
但如果我收到一条带有key = S9 的流消息,则不会发生加入,然后我想记录一条消息:
------- join failed, key not found ---------
据我所知,这可以使用filter 方法来实现,但我不确定如何实现。请帮助我如何做到这一点,或者有没有更好的方法来做到这一点。
【问题讨论】:
-
你正在做一个左连接,它总是会成功,你将从左数据帧中获取所有数据并从右数据帧中获取匹配数据。你能用一些示例数据和你期望的输出来更新这个问题吗?
-
嘿@NikunjKakadiya,我已经更新了问题。
-
您能否指定您的两个数据框内容。您添加的内容对获得想法没有多大帮助
-
@NikunjKakadiya,已更新。
-
不知何故,我需要检查消息(丰富后)是否有额外的列
site。如果它不存在,那么我需要记录一个语句。
标签: scala apache-spark apache-spark-sql apache-kafka-streams