【问题标题】:2 spark stream job with same consumer group id2 具有相同消费者组 ID 的火花流作业
【发布时间】:2018-11-06 04:17:08
【问题描述】:

我正在尝试对消费者群体进行实验

这是我的代码 sn-p

public final class App {

private static final int INTERVAL = 5000;

public static void main(String[] args) throws Exception {

    Map<String, Object> kafkaParams = new HashMap<>();
    kafkaParams.put("bootstrap.servers", "xxx:9092");
    kafkaParams.put("key.deserializer", StringDeserializer.class);
    kafkaParams.put("value.deserializer", StringDeserializer.class);
    kafkaParams.put("auto.offset.reset", "earliest");
    kafkaParams.put("enable.auto.commit", true);
    kafkaParams.put("auto.commit.interval.ms","1000");
    kafkaParams.put("security.protocol","SASL_PLAINTEXT");
    kafkaParams.put("sasl.kerberos.service.name","kafka");
    kafkaParams.put("retries","3");
    kafkaParams.put(GROUP_ID_CONFIG,"mygroup");
    kafkaParams.put("request.timeout.ms","210000");
    kafkaParams.put("session.timeout.ms","180000");
    kafkaParams.put("heartbeat.interval.ms","3000");
    Collection<String> topics = Arrays.asList("venkat4");

    SparkConf conf = new SparkConf();
    JavaStreamingContext ssc = new JavaStreamingContext(conf, new Duration(INTERVAL));


    final JavaInputDStream<ConsumerRecord<String, String>> stream =
            KafkaUtils.createDirectStream(
                    ssc,
                    LocationStrategies.PreferConsistent(),
                    ConsumerStrategies.<String, String>Subscribe(topics, kafkaParams)
            );

    stream.mapToPair(
            new PairFunction<ConsumerRecord<String, String>, String, String>() {
                @Override
                public Tuple2<String, String> call(ConsumerRecord<String, String> record) {
                    return new Tuple2<>(record.key(), record.value());
                }
            }).print();


    ssc.start();
    ssc.awaitTermination();


}

}

当我同时运行两个 Spark 流式传输作业时,它会因错误而失败

线程“main”java.lang.IllegalStateException 中的异常:分区 venkat4-1 没有当前分配 在 org.apache.kafka.clients.consumer.internals.SubscriptionState.assignedState(SubscriptionState.java:251) 在 org.apache.kafka.clients.consumer.internals.SubscriptionState.needOffsetReset(SubscriptionState.java:315) 在 org.apache.kafka.clients.consumer.KafkaConsumer.seekToEnd(KafkaConsumer.java:1170) 在 org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.latestOffsets(DirectKafkaInputDStream.scala:197) 在 org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.compute(DirectKafkaInputDStream.scala:214) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:341) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:341) 在 scala.util.DynamicVariable.withValue(DynamicVariable.scala:58) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:340) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:340) 在 org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:415) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:335) 在 org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:333) 在 scala.Option.orElse(Option.scala:289)

根据 https://www.wisdomjobs.com/e-university/apache-kafka-tutorial-1342/apache-kafka-consumer-group-example-19004.html 创建具有相同组的单独的 kafka 消费者实例将创建分区的重新平衡。我相信消费者不会容忍再平衡。我该如何解决这个问题

下面是使用的命令

SPARK_KAFKA_VERSION=0.10 spark2-submit --num-executors 2 --master yarn --deploy-mode client --files jaas.conf#jaas.conf,hive.keytab#hive.keytab --driver-java-options "-Djava.security.auth.login.config=./jaas.conf" --class Streaming.App --conf "spark.executor.extraJavaOptions=-Djava.security.auth.login.config=./jaas.conf " --conf spark.streaming.kafka.consumer.cache.enabled=false 1-1.0-SNAPSHOT.jar

【问题讨论】:

  • 到目前为止,我尝试配置以下参数,但这些参数都不起作用 kafkaParams.put("request.timeout.ms","210000"); kafkaParams.put("session.timeout.ms","180000"); kafkaParams.put("heartbeat.interval.ms","3000"); kafkaParams.put("metadata.max.age.ms","1000"); spark.streaming.kafka.consumer.cache.enabled=false

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


【解决方案1】:

@Ravikumar 为延迟道歉。

我的测试是这样完成的

一个。我的主题有 3 个分区 湾。火花流作业由 2 个执行程序启动——运行良好。 C。后来我决定用另一个实例来扩展它,方法是运行另一个带有 1 个执行器的 spark-streaming 作业,以匹配我失败的第三个分区。

关于您的陈述: 当您启动第二个 spark 流作业时,另一个消费者尝试使用来自同一个消费者 groupid 的同一个分区。所以它抛出错误 是的,这完全正确。但为什么它不容忍是这里的问题。

引用您突出显示的文档:

Kafka 将主题的分区分配给组中的消费者,因此 每个分区只被组中的一个消费者消费。 Kafka 保证一条消息只能被单个消费者读取 在群里。 Kafka 随时重新平衡分区存储 代理失败或向现有主题添加新分区。这是 kafka 具体如何在分区中平衡数据 经纪人。如果添加更多进程/线程,Kafka 将重新平衡。 ZooKeeper 可以由 Kafka 集群重新配置,如果有任何消费者或 代理无法向 ZooKeeper 发送心跳。

这也是我对火花流工作的期望。我尝试了能够容忍重新平衡的普通 kafka 客户端。

您在文档“缓存由 topicpartition 和 group.id 键控,因此每次调用 createDirectStream 使用单独的 group.id”中的观点澄清了我的问题。

另外来自 PR https://github.com/apache/spark/pull/21038 -- 以下是提及

"当新的消费者加入时,Kafka 分区可以被撤销 消费者组来重新平衡分区。但目前的 Spark Kafka 连接器代码确保没有分区撤销场景,所以 试图从撤销的分区中获取最新的偏移量会抛出 JIRA 提到的例外情况。”

很高兴关闭此线程。非常感谢您的回复

【讨论】:

    【解决方案2】:

    根据 https://www.wisdomjobs.com/e-university/apache-kafka-tutorial-1342/apache-kafka-consumer-group-example-19004.html 创建具有相同组的单独的 kafka 消费者实例将创建分区的重新平衡。我相信消费者不会容忍再平衡。我该如何解决这个问题

    现在所有分区都只被一个消费者消费。如果数据摄取率很高,消费者可能会以摄取的速度消费数据。

    在同一个consumergroup中增加更多的consumer,消费一个topic的数据,提高消费率。使用这种方法的 Spark 流在 Kafka 分区和 Spark 分区之间实现 1:1 并行性。 Spark 将在内部处理它。

    如果消费者数量多于主题分区,它将处于空闲状态并且资源未得到充分利用。始终建议消费者应小于或等于分区数。

    如果添加更多进程/线程,Kafka 将重新平衡。如果任何消费者或代理无法向 ZooKeeper 发送心跳,则 ZooKeeper 可以由 Kafka 集群重新配置。

    每当任何代理失败或向现有主题添加新分区时,Kafka 都会重新平衡分区存储。这是 kafka 特定的如何在代理中跨分区平衡数据。

    Spark 流式处理在 Kafka 分区和 Spark 分区之间提供简单的 1:1 并行性。如果您没有使用 ConsumerStragies.Assign 提供任何分区详细信息,则使用给定主题的所有分区。

    Kafka 将主题的分区分配给组中的消费者,因此 每个分区只被组中的一个消费者消费。 Kafka 保证一条消息只能被单个消费者读取 在群里。

    当您启动第二个 spark 流作业时,另一个消费者尝试使用来自同一个消费者 groupid 的同一个分区。所以它会抛出错误。

    val alertTopics = Array("testtopic")
    
    val kafkaParams = Map[String, Object](
      "bootstrap.servers" -> sparkJobConfig.kafkaBrokers,
      "key.deserializer" -> classOf[StringDeserializer],
      "value.deserializer" -> classOf[StringDeserializer],
      "group.id" -> sparkJobConfig.kafkaConsumerGroup,
      "auto.offset.reset" -> "latest"
    )
    
    val streamContext = new StreamingContext(sparkContext, Seconds(sparkJobConfig.streamBatchInterval.toLong))
    
    val streamData = KafkaUtils.createDirectStream(streamContext, PreferConsistent, Subscribe[String, String](alertTopics, kafkaParams))
    

    如果要使用特定于分区的 spark 作业,请使用以下代码。

    val topicPartitionsList =  List(new TopicPartition("topic",1))
    
    val alertReqStream1 = KafkaUtils.createDirectStream(streamContext, PreferConsistent, ConsumerStrategies.Assign(topicPartitionsList, kafkaParams))
    

    https://spark.apache.org/docs/2.2.0/streaming-kafka-0-10-integration.html#consumerstrategies

    消费者可以使用相同的group.id加入群组。

    val topicPartitionsList =  List(new TopicPartition("topic",3), new TopicPartition("topic",4))
    
        val alertReqStream2 = KafkaUtils.createDirectStream(streamContext, PreferConsistent, ConsumerStrategies.Assign(topicPartitionsList, kafkaParams))
    

    再添加两个消费者就是添加到同一个 groupid 中。

    请阅读 Spark-Kafka 集成指南。 https://spark.apache.org/docs/2.2.0/streaming-kafka-0-10-integration.html

    希望这会有所帮助。

    【讨论】:

    • @venkat-sam 能否请您确认一下,如果这澄清了问题,它将对其他用户有用。
    猜你喜欢
    • 1970-01-01
    • 2021-12-04
    • 2020-11-21
    • 2023-03-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-03-19
    • 2020-04-19
    相关资源
    最近更新 更多