【问题标题】:Watermarking for Spark structured streaming with three way joins使用三向连接的 Spark 结构化流的水印
【发布时间】:2019-09-26 14:49:25
【问题描述】:

我有 3 个数据流:foobarbaz

有必要将这些流与LEFT OUTER JOIN 连接到以下链中:foo -> bar -> baz

这里尝试使用内置的rate 流来模拟这些流:

val rateStream = session.readStream
  .format("rate")
  .option("rowsPerSecond", 5)
  .option("numPartitions", 1)
  .load()

val fooStream = rateStream
  .select(col("value").as("fooId"), col("timestamp").as("fooTime"))

val barStream = rateStream
  .where(rand() < 0.5) // Introduce misses for ease of debugging
  .select(col("value").as("barId"), col("timestamp").as("barTime"))

val bazStream = rateStream
  .where(rand() < 0.5) // Introduce misses for ease of debugging
  .select(col("value").as("bazId"), col("timestamp").as("bazTime"))

这是将所有这些流连接在一起的第一种方法,假设foobarbaz 的潜在延迟很小(~5 seconds):

val foobarStream = fooStream
  .withWatermark("fooTime", "5 seconds")
  .join(
    barStream.withWatermark("barTime", "5 seconds"),
    expr("""
       barId = fooId AND
       fooTime >= barTime AND
       fooTime <= barTime + interval 5 seconds
           """),
    joinType = "leftOuter"
  )

val foobarbazQuery = foobarStream
  .join(
    bazStream.withWatermark("bazTime", "5 seconds"),
    expr("""
      bazId = fooId AND
      bazTime >= fooTime AND
      bazTime <= fooTime + interval 5 seconds
         """),
    joinType = "leftOuter")
  .writeStream
  .format("console")
  .start()

通过上面的设置,我可以观察到以下数据元组:

  • (some_foo, some_bar, some_baz)
  • (some_foo, some_bar, null)

但仍然缺少(some_foo, null, some_baz)(some_foo, null, null)

任何想法,如何正确配置水印以获得所有组合?

更新:

barTime 上为foobarStream 添加了额外的水印之后:

val foobarbazQuery = foobarStream
  .withWatermark("barTime", "1 minute")
  .join(/* ... */)`

我能够得到这个(some_foo, null, some_baz) 组合,但仍然缺少(some_foo, null, null)...

【问题讨论】:

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


    【解决方案1】:

    我留下一些信息仅供参考。

    链式流-流连接无法正常工作,因为 Spark 仅支持全局水印(而不是操作员水印),这可能会导致连接之间的中间输出丢失。

    Apache Spark 社区指出了这个问题,并在之前讨论过。以下是更多详细信息的链接: https://lists.apache.org/thread.html/cc6489a19316e7382661d305fabd8c21915e5faf6a928b4869ac2b4a@%3Cdev.spark.apache.org%3E

    (免责声明:我是发起邮件线程的作者。)

    【讨论】:

    • 感谢分享!我希望这个问题早点被发现,这样我在 2018 年夏天就不会那么挣扎了 :)
    猜你喜欢
    • 2019-04-06
    • 2021-05-03
    • 2021-05-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多