【发布时间】:2020-11-04 13:43:28
【问题描述】:
我有一个场景,我想将相同的流数据帧保存到两个不同的流接收器。
我创建了一个流式数据帧,我需要将其发送到 Kafka 主题和 delta Lake。
我想过使用 forEachBatch,但看起来它不支持多个 STREAMING SINKS。
另外,我尝试将 spark session.awaitAnyTermination() 与多个写入流一起使用。但是第二个流没有得到处理!
有什么方法可以实现吗?!
这是我的代码:
- 我正在从 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)]
- 将上述数据帧写入 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()
- 将相同的流式数据帧写入 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