【问题标题】:Spark Streaming processing from multiple rabbitmq queue in parallel来自多个rabbitmq队列的Spark Streaming并行处理
【发布时间】:2017-02-26 03:31:40
【问题描述】:

我试图为多个 RabbitMQ 队列设置 Spark 流。如下所述,我设置了 2 个工人,每个工人都有一个核心和 2GB 内存。所以,问题是当我将此参数保持为conf.set("spark.cores.max","2") 时,流不处理任何数据,它只是继续添加作业。但是一旦我将它设置为conf.set("spark.cores.max","3"),流式传输就会开始处理它。所以,我无法理解这样做的原因。另外,如果我想从两个队列中并行处理数据,我应该怎么做。我在下面提到了我的代码和配置设置。

Spark-env.sh:

SPARK_WORKER_MEMORY=2g 
SPARK_WORKER_INSTANCES=1
SPARK_WORKER_CORES=1

Scala 代码:

    val rabbitParams =  Map("storageLevel" -> "MEMORY_AND_DISK_SER_2","queueName" -> config.getString("queueName"),"host" -> config.getString("QueueHost"), "exchangeName" -> config.getString("exchangeName"), "routingKeys" -> config.getString("routingKeys"))
    val receiverStream = RabbitMQUtils.createStream(ssc, rabbitParams)
    receiverStream.start()      

    val predRabbitParams =  Map("storageLevel" -> "MEMORY_AND_DISK_SER_2", "queueName" -> config.getString("queueName1"), "host" -> config.getString("QueueHost"), "exchangeName" -> config.getString("exchangeName1"), "routingKeys" -> config.getString("routingKeys1"))
    val predReceiverStream = RabbitMQUtils.createStream(ssc, predRabbitParams)
    predReceiverStream.start()  

【问题讨论】:

    标签: apache-spark spark-streaming datastax


    【解决方案1】:

    Streaming Guide 中解释了此行为。每个接收器都是一个长时间运行的进程,它占用一个线程。

    如果可用线程的数量小于或等于接收者的数量,则没有资源用于任务处理:

    分配给 Spark Streaming 应用程序的核心数量必须大于接收器的数量。否则系统会收到数据,但无法处理。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-05-02
      • 1970-01-01
      • 2017-12-18
      • 2017-03-29
      • 1970-01-01
      • 1970-01-01
      • 2022-12-09
      • 1970-01-01
      相关资源
      最近更新 更多