【问题标题】:Spark Structure streaming read data twice per every micro-batch. How to avoidSpark Structure 每个微批次流式读取数据两次。如何避免
【发布时间】:2020-07-22 22:50:12
【问题描述】:

我对 spark 结构流有一个非常奇怪的问题。 Spark 结构流式处理为每个微批次创建两个 Spark 作业。 结果,从 Kafka 读取数据两次。 这是一个简单的代码sn-p。

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.Trigger

object CheckHowSparkReadFromKafka {
  def main(args: Array[String]): Unit = {
    val session = SparkSession.builder()
      .config(new SparkConf()
        .setAppName(s"simple read from kafka with repartition")
        .setMaster("local[*]")
        .set("spark.driver.host", "localhost"))
      .getOrCreate()
    val testPath = "/tmp/spark-test"
    FileSystem.get(session.sparkContext.hadoopConfiguration).delete(new Path(testPath), true)
    import session.implicits._
    val stream = session
      .readStream
      .format("kafka")
      .option("kafka.bootstrap.servers",        "kafka-20002-prod:9092")
      .option("subscribe", "topic")
      .option("maxOffsetsPerTrigger", 1000)
      .option("failOnDataLoss", false)
      .option("startingOffsets", "latest")
      .load()
      .repartitionByRange( $"offset")
      .writeStream
      .option("path", testPath + "/data")
      .option("checkpointLocation", testPath + "/checkpoint")
      .format("parquet")
      .trigger(Trigger.ProcessingTime(10.seconds))
      .start()
    stream.processAllAvailable()

发生这种情况是因为如果.repartitionByRange( $"offset"),如果我删除此行,一切都很好。 但是使用 spark 创建两个作业,一个具有 1 个阶段刚刚从 Kafka 读取,第二个具有 3 个阶段读取 -> 洗牌 -> 写入。 所以第一份工作的结果从未使用过。

这会对性能产生重大影响。 我的一些 Kafka 主题有 1550 个分区,因此阅读它们两次很重要。 如果我添加缓存,事情会变得更好,但这对我来说不是一种方式。 在本地模式下,批处理中的第一个作业用时不到 0.1 毫秒,但索引为 0 的批处理除外。但在 YARN 集群和 Messos 中,这两个作业完全符合预期,而我的主题则需要将近 1.2 分钟。

为什么会这样?我怎样才能避免这种情况?看起来像虫子?

附:我使用火花 2.4.3。

【问题讨论】:

  • 你能看到write阶段正在创建的任务号和shuffle write/read吗?
  • @EmiCareOfCell44 是的,我可以看到所有任务编号和阶段。第一个 Job 有 1 个阶段,其中任务数为 240(与 Kafka 中的分区数相同)。第二个 Job 有 2 个 stage,其中一个 stage 类似于第一个 Job,第二个 stage 是 shuffle + write。
  • 我想对于范围分区器,Spark 必须首先读取所有消息才能创建 RangePartiitoner。因为偏移量是每个分区中的一个未知数,它读取并重新分配一个分区中的所有消息以创建de索引,然后对每个执行器要处理的数据进行混洗。您的案例必须使用范围分区器吗?
  • 1.是的,RangePartitioner 是强制性的。 2. Spark 不能在另一个作业中重用作业的结果,除非将其持久化到某个接收器并再次读取。并且大多数和第二个工作1阶段的事情和第一个阶段的唯一阶段完全一样。更重要的是,如果我将使用直接流式传输并从每个微批次红色创建一个数据帧,然后使用相同的逻辑 - 我将有 1 个作业。
  • 当我使用简单排序时也会发生同样的情况。

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


【解决方案1】:

在这种情况下,火花中没有错误。 从 Kafka 读取此数据两次的根本原因非常简单。 repartitionByRange 函数生成两个 Spark 作业。

一个用于实际重新分区。

一个用于采样以查找分区的边界。

更多详情请到spark jira

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2011-06-06
    • 1970-01-01
    • 2020-08-06
    • 1970-01-01
    • 2015-07-14
    • 1970-01-01
    • 2021-11-17
    • 2019-10-06
    相关资源
    最近更新 更多