【问题标题】:DSE Spark Streaming: Long active batches queueDSE Spark Streaming:长活动批次队列
【发布时间】:2018-10-04 16:55:50
【问题描述】:

我有以下代码:

val conf = new SparkConf()
  .setAppName("KafkaReceiver")
  .set("spark.cassandra.connection.host", "192.168.0.78")
  .set("spark.cassandra.connection.keep_alive_ms", "20000")
  .set("spark.executor.memory", "2g")
  .set("spark.driver.memory", "4g")
  .set("spark.submit.deployMode", "cluster")
  .set("spark.executor.instances", "3")
  .set("spark.executor.cores", "3")
  .set("spark.shuffle.service.enabled", "false")
  .set("spark.dynamicAllocation.enabled", "false")
  .set("spark.io.compression.codec", "snappy")
  .set("spark.rdd.compress", "true")
  .set("spark.streaming.backpressure.enabled", "true")
  .set("spark.streaming.backpressure.initialRate", "200")
  .set("spark.streaming.receiver.maxRate", "500")

val sc = SparkContext.getOrCreate(conf)
val ssc = new StreamingContext(sc, Seconds(10))
val sqlContext = new SQLContext(sc)
val kafkaParams = Map[String, String](
  "bootstrap.servers" -> "192.168.0.113:9092",
  "group.id" -> "test-group-aditya",
  "auto.offset.reset" -> "largest")

val topics = Set("random")
val kafkaStream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc, kafkaParams, topics)

我正在使用以下命令通过spark-submit 运行代码:

dse> bin/dse spark-submit --class test.kafkatesting /home/aditya/test.jar

我在不同的机器上安装了一个三节点 Cassandra DSE 集群。每当我运行应用程序时,它都会占用大量数据并开始创建活动批次队列,这反过来又会产生积压和长时间的调度延迟。如何提高性能并控制队列,使其仅在执行完当前批次后才接收新批次?

【问题讨论】:

    标签: scala apache-spark spark-streaming cassandra-3.0 spark-cassandra-connector


    【解决方案1】:

    我找到了解决方案,在代码中做了一些优化。与其保存 RDD,不如尝试创建 Dataframe,与 RDD 相比,将 DF 保存到 Cassandra 的速度要快得多。另外,为了获得好的效果,增加core和executor内存的数量。

    谢谢,

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-10-15
      • 2018-12-26
      • 1970-01-01
      • 1970-01-01
      • 2015-05-21
      • 2021-07-13
      • 1970-01-01
      • 2016-05-04
      相关资源
      最近更新 更多