【问题标题】:Spark structured streaming - UNION two or more streaming sourcesSpark 结构化流 - UNION 两个或多个流源
【发布时间】: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


【解决方案1】:

首先,请定义您的案例类 OffsetSeq 与数据帧联合的代码之间的关系。

接下来,在执行此联合然后使用 writestream 写入 Kafka 时,检查点是一个真正的问题。分成多个写流——每个都有自己的检查点——因为联合操作而混淆了批处理 id。由于检查点似乎在寻找在联合之前生成数据帧的所有模型,并且无法区分哪些行/记录来自哪个数据帧/模型,因此使用相同的写入流与数据帧的联合会失败并设置检查点。

对于从结构化 sql 流联合数据帧写入 Kafka - 最好将 writestream 与 foreach 和 ForEachWriter 一起使用,包括流程方法中的 Kafka Producer。不需要检查点;应用程序仅使用临时检查点文件,这些文件设置为在适当时删除 - 在会话生成器中将“forceDeleteTempCheckpointLocation”设置为 true。

无论如何,我刚刚设置了 scala 代码来合并任意数量的流数据帧,然后写入 Kafka Producer。将所有 Kafka Producer 代码放入 ForEachWriter 进程方法后,似乎可以正常工作,以便 Spark 可以对其进行序列化。

val output = dataFrameModelArray.reduce(_ union _)
val stream: StreamingQuery = output
  .writeStream.foreach(new ForeachWriter[Row] {

    def open(partitionId: Long, version: Long): Boolean = {
      true
    }

    def process(row: Row): Unit = {
      val producer: KafkaProducer[String, String] = new KafkaProducer[String, String](props)
      val record = new ProducerRecord[String, String](producerTopic, row.getString(0), row.getString(1))
      producer.send(record)
    }

    def close(errorOrNull: Throwable): Unit = {
    }
  }
).start()

如果需要,可以在处理方法中添加更多逻辑。

注意在联合之前,所有要联合的数据框都已转换为键、值字符串列。值是要通过 Kafka Producer 发送的消息数据的 json 字符串。这对于在尝试联合之前写入也非常重要。

svcModel.transform(query)
    .select($"key", $"uuid", $"currentTime", $"label", $"rawPrediction", $"prediction")
    .selectExpr("key", "to_json(struct(*)) AS value")
    .selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")

其中 svcModel 是 dataFrameModelArray 中的一个数据框。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-09-16
    • 2019-12-13
    • 1970-01-01
    • 1970-01-01
    • 2017-05-04
    • 2018-11-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多