【问题标题】:Using two WriteStreams in same spark structured streaming job在同一个 Spark 结构化流作业中使用两个 WriteStream
【发布时间】:2020-11-04 13:43:28
【问题描述】:

我有一个场景,我想将相同的流数据帧保存到两个不同的流接收器

我创建了一个流式数据帧,我需要将其发送到 Kafka 主题和 delta Lake。

我想过使用 forEachBatch,但看起来它不支持多个 STREAMING SINKS。

另外,我尝试将 spark session.awaitAnyTermination() 与多个写入流一起使用。但是第二个流没有得到处理!

有什么方法可以实现吗?!

这是我的代码:

  1. 我正在从 Kafka 流中读取数据并创建单个流数据帧。
val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9092") 
  .option("subscribe", "ingestionTopic1")
  .load() 

df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)").as[(String, String)]
  1. 将上述数据帧写入 Kafka 主题
val ds1 = df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
  .writeStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "localhost:9082")
  .option("topic", "outputTopic1")
  .start()
  1. 将相同的流式数据帧写入 delta Lake
val ds2 = df.format("delta")
  .outputMode("append")
  .option("checkpointLocation", "/test/delta/events/_checkpoints/etlflow") 
  .start("/test/delta/events") 

ds1.awaitTermination 
ds2.awaitTermination

【问题讨论】:

  • This 为我工作,也许它也解决了你的问题。
  • 感谢迈克的分享。我查看了该代码 sn-p,我们正在使用两个不同的流数据帧。就我而言,我使用的是来自 Kafka 的单个流数据帧。如果我使用该数据帧启动一个 WriteStream() 操作,则偏移量被消耗并已移动到 LATEST 并且第二个 writeStream() 操作没有消耗它,因为它已被第一个消耗。
  • 嗨@mike - 非常感谢您的编辑。我尝试为每个 writeStreams() 使用不同的检查点位置,并在启动两个查询后添加了 awaitTermination()。但我仍然只能看到数据只写入一个接收器而不是另一个接收器。
  • 我还查看了每个接收器的偏移量,只有一个接收器的偏移量在进行中,而另一个接收器的偏移量保持不变并且没有进行。
  • 您是否在第二个查询中添加了writeStream

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


【解决方案1】:

要将一个输入流用于多个输出流,您需要遵循以下几点:

  • 您需要确保在两个输出流中有两个不同的检查点位置。

  • 此外,您需要确保在第二个输出查询中也有 writeStream 调用。

  • 总体而言,在等待两个查询终止之前启动这两个查询非常重要。 (你已经这样做了)

【讨论】:

  • 检查点位置提示她很有用,因为(至少对我而言)这是我面临的多个流查询问题的不直观原因
猜你喜欢
  • 1970-01-01
  • 2019-11-12
  • 1970-01-01
  • 2018-12-12
  • 1970-01-01
  • 2017-09-16
  • 2019-09-10
  • 2019-10-07
  • 2021-09-20
相关资源
最近更新 更多