【问题标题】:Spark Structured Streaming Multiple WriteStreams to Same SinkSpark结构化流多个WriteStreams到同一个接收器
【发布时间】:2018-12-24 19:21:49
【问题描述】:

在 Spark Structured Streaming 2.2.1 中,同一个数据库接收器的两个 Writestream 不会按顺序发生。请建议如何让它们按顺序执行。

val deleteSink = ds1.writestream
  .outputMode("update")
  .foreach(mydbsink)
  .start()

val UpsertSink = ds2.writestream
  .outputMode("update")
  .foreach(mydbsink)
  .start()

deleteSink.awaitTermination()
UpsertSink.awaitTermination()

使用上面的代码,deleteSinkUpsertSink 之后执行。

【问题讨论】:

    标签: scala apache-spark slick-3.0 spark-structured-streaming


    【解决方案1】:

    如果你想让两个流并行运行,你必须使用

    sparkSession.streams.awaitAnyTermination()
    

    而不是

    deleteSink.awaitTermination()
    UpsertSink.awaitTermination()
    

    在您的情况下,除非 deleteSink 将停止或抛出异常,否则 UpsertSink 将永远不会启动,正如它在 scaladoc 中所说的那样

    等待this 查询终止,无论是query.stop() 还是异常。 如果查询因异常而终止,则将引发异常。 如果查询已终止,则对该方法的所有后续调用都将返回 立即(如果查询被stop() 终止),或者抛出异常 立即(如果查询因异常而终止)。

    【讨论】:

    • 我的问题是对同一个 dbsink 的两个写入流并不总是连续的。有时我看到 upsert sink 在 delete sink 之前被执行。
    • 嗯,在流式传输的上下文中,想要有一个序列并没有什么意义,因为过了一段时间,你无法确定接下来会执行哪个查询。也许你应该考虑改变你处理这个问题的方式。例如,如果两个 DS 具有相同的架构,您可以执行 ds1.union(ds2) 并修改 'mydbsink' 以检查某个值(例如,预先添加一个可以有两个值的列 Action:删除或插入)。然后就可以group by,key排序,然后使用修改后的dbsink。也许你有这样更大的控制权。
    • 是的,我使用了类似的方法来解决问题。
    • 如果一个写入流输出是另一个写入流的输入,是否有办法让写入流始终按顺序排列?
    • 嗨,Alex,我们只使用了一个 Sink 而不是两个不同的 sink 用于 delete 和 Upsert 。我已经根据唯一键对记录进行了重新分区,并在 slick 中使用可变集来检查和删除 DB 中的记录,如果该集不包含键而与操作指示符无关,并且还基于操作指示符,我从 Postgresql DB 中插入或删除了记录。
    猜你喜欢
    • 2018-04-19
    • 2020-09-06
    • 1970-01-01
    • 2019-11-12
    • 2019-09-21
    • 1970-01-01
    • 1970-01-01
    • 2020-10-27
    • 1970-01-01
    相关资源
    最近更新 更多