【发布时间】:2019-11-12 11:27:12
【问题描述】:
我正在使用 spark 2.3.2 并遇到一个问题,即对来自 Kafka 的 2 个或更多流媒体源进行联合。其中每一个都是来自 Kafka 的流式源,我已经转换并存储在 Dataframes 中。
理想情况下,我希望将这个 UNIONed 数据帧的结果以 parquet 格式存储在 HDFS 中,甚至可能存储回 Kafka。最终目标是以尽可能低的延迟存储这些合并的事件。
val finalDF = flatDF1
.union(flatDF2)
.union(flatDF3)
val query = finalDF.writeStream
.format("parquet")
.outputMode("append")
.option("path", hdfsLocation)
.option("checkpointLocation", checkpointLocation)
.option("failOnDataLoss", false)
.start()
query.awaitTermination()
当对控制台而不是 parquet 执行 writeStream 时,我得到了预期的结果,但上面的示例导致断言失败。
Caused by: java.lang.AssertionError: assertion failed
at scala.Predef$.assert(Predef.scala:156)
at org.apache.spark.sql.execution.streaming.OffsetSeq.toStreamProgress(OffsetSeq.scala:42)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution.org$apache$spark$sql$execution$streaming$MicroBatchExecution$$populateStartOffsets(MicroBatchExecution.scala:185)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply$mcV$sp(MicroBatchExecution.scala:124)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply(MicroBatchExecution.scala:121)
at org.apache.spark.sql.execution.streaming.MicroBatchExecution$$anonfun$runActivatedStream$1$$anonfun$apply$mcZ$sp$1.apply(MicroBatchExecution.scala:121)
这是失败的类和断言:
case class OffsetSeq(offsets: Seq[Option[Offset]], metadata: Option[OffsetSeqMetadata] = None) {
assert(sources.size == offsets.size)
这是因为检查点仅存储其中一个数据帧的偏移量吗?查看 Spark Structured Streaming 文档,似乎可以在 Spark 2.2 或 >
中进行流源的连接/联合【问题讨论】:
-
你为什么使用检查点而不是手动提交?
-
@maximeG 是否可以在没有检查点的情况下进行结构化流式传输?你还能如何维护已经从 Kafka 消费的内容的状态?
-
kafka 本身为您的消费者提供了一个主题:__consumer_offsets
-
您应该查看 spark/kafka 文档以进行手动提交:spark.apache.org/docs/latest/… 与检查点相比,如果您遇到 kafka 问题,无论您的应用程序代码如何更改,Kafka 都是一个持久存储,就像您两次读取相同的数据一样,请查看您的 kafka 消费者中的 group.id 和 auto.offset.reset 属性
标签: scala apache-spark union spark-structured-streaming