代码之家  ›  专栏  ›  技术社区  ›  Jorge Cespedes

Spark流媒体1.6+Kafka:处于“排队”状态的批太多

  •  1
  • Jorge Cespedes  · 技术社区  · 8 年前

    我使用spark流来消费来自一个kafka主题的消息,这个主题有10个分区。我使用的是直接从卡夫卡消费的方法,代码如下:

    def createStreamingContext(conf: Conf): StreamingContext = {
        val dateFormat = conf.dateFormat.apply
        val hiveTable = conf.tableName.apply
    
        val sparkConf = new SparkConf()
    
        sparkConf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
        sparkConf.set("spark.driver.allowMultipleContexts", "true")
    
        val sc = SparkContextBuilder.build(Some(sparkConf))
        val ssc = new StreamingContext(sc, Seconds(conf.batchInterval.apply))
    
        val kafkaParams = Map[String, String](
          "bootstrap.servers" -> conf.kafkaBrokers.apply,
          "key.deserializer" -> classOf[StringDeserializer].getName,
          "value.deserializer" -> classOf[StringDeserializer].getName,
          "auto.offset.reset" -> "smallest",
          "enable.auto.commit" -> "false"
        )
    
        val directKafkaStream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](
          ssc,
          kafkaParams,
          conf.topics.apply().split(",").toSet[String]
        )
    
        val windowedKafkaStream = directKafkaStream.window(Seconds(conf.windowDuration.apply))
        ssc.checkpoint(conf.sparkCheckpointDir.apply)
    
        val eirRDD: DStream[Row] = windowedKafkaStream.map { kv =>
          val fields: Array[String] = kv._2.split(",")
          createDomainObject(fields, dateFormat)
        }
    
        eirRDD.foreachRDD { rdd =>
          val schema = SchemaBuilder.build()
          val sqlContext: HiveContext = HiveSQLContext.getInstance(Some(rdd.context))
          val eirDF: DataFrame = sqlContext.createDataFrame(rdd, schema)
    
          eirDF
            .select(schema.map(c => col(c.name)): _*)
            .write
            .mode(SaveMode.Append)
            .partitionBy("year", "month", "day")
            .insertInto(hiveTable)
        }
        ssc
      }
    

    从代码中可以看出,我使用window来实现( 如果我错了请纠正我 ):因为有一个操作要插入到配置单元表中,所以我希望避免过于频繁地写入hdfs,所以我希望在内存中保留足够的数据,然后才写入文件系统。我认为使用窗口是实现这一目标的正确方法。

    现在,在下面的图像中,您可以看到有许多批正在排队,并且正在处理的批将永远完成。

    Only one batch being in processing status and others queued forever

    我还提供了正在处理的单个批次的详细信息:

    Insert into generates thousands of tasks!! Why

    当批处理中没有太多事件时,为什么insert操作有这么多任务?有时,拥有0个事件也会生成数千个任务,这些任务需要永远才能完成。

    我处理带有spark的微区的方法是错误的吗?

    谢谢你的帮助!

    一些额外的细节:

    纱线容器的最大容量为2gb。 在这个纱线队列中,容器的最大数量是10个。 当我查看执行这个spark应用程序的队列的详细信息时,容器的数量非常大,大约有15k个挂起的容器。

    2 回复  |  直到 8 年前
        1
  •  0
  •   Jorge Cespedes    8 年前

    好吧,我终于明白了。显然spark流不能处理空事件,所以在代码的foreachrdd部分中,我添加了以下内容:

    eirRDD.foreachRDD { rdd =>
          if (rdd.take(1).length != 0) {
            //do action
          }
    }
    

    这样我们就跳过了空的微批量。isEmpty()方法不起作用。

    希望这能帮助别人!;)

        2
  •  0
  •   lvnt    8 年前

    您可以使用rdd.partitions.size进行检查。如果大小小于1,则表示RDD为空。但如果使用take()方法,则在rdd为空时会出错。

    推荐文章