You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

同Consumer Group ID的两个Spark Streaming任务并发运行报错求助

解决同Consumer Group下Spark Streaming任务启动的分区重平衡报错问题

我来帮你搞定这个问题!你遇到的IllegalStateException: No current assignment for partition venkat4-1错误,本质是Spark Streaming的DirectKafkaInputDStream在尝试获取最新偏移量时,Kafka的分区重平衡还没完成,导致消费者还没拿到分区分配信息就执行了seek操作。下面一步步拆解解决方案:

问题场景回顾

你正在做Kafka Consumer Group实验,运行两个使用相同Group ID(mygroup)的Spark Streaming任务,任务代码和提交命令如下:

你的Spark Streaming代码片段

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(); 
    } 
}

报错日志

Exception in thread "main" java.lang.IllegalStateException: No current assignment for partition venkat4-1
at org.apache.kafka.clients.consumer.internals.SubscriptionState.assignedState(SubscriptionState.java:251)
at org.apache.kafka.clients.consumer.internals.SubscriptionState.needOffsetReset(SubscriptionState.java:315)
at org.apache.kafka.clients.consumer.KafkaConsumer.seekToEnd(KafkaConsumer.java:1170)
at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.latestOffsets(DirectKafkaInputDStream.scala:197)
at org.apache.spark.streaming.kafka010.DirectKafkaInputDStream.compute(DirectKafkaInputDStream.scala:214)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:341)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:341)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:58)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:340)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:340)
at org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:415)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:335)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:333)
at scala.Option.orElse(Option.scala:289)

核心问题分析

当同Group的第二个消费者启动时,Kafka会触发分区重平衡,但你的消费者配置和Spark Streaming的处理逻辑没适配这个过程:

  1. heartbeat.interval.ms设置过小(3000ms),和session.timeout.ms(180000ms)的比例严重偏离Kafka推荐的1/3原则,可能导致协调器误判消费者状态
  2. 自动提交偏移量(enable.auto.commit=true)在重平衡时容易出现偏移量提交冲突
  3. DirectKafkaInputDStream在重平衡完成前就尝试获取偏移量,导致找不到分区分配信息

具体解决方案

1. 调整Kafka消费者核心参数

修改kafkaParams中的以下配置,适配重平衡流程:

  • 关闭自动提交偏移量,改为手动提交
  • 增大max.poll.interval.ms,给重平衡足够的时间窗口
  • 修正heartbeat.interval.ms和session.timeout.ms的比例(推荐heartbeat为session timeout的1/3)

修改后的kafkaParams关键部分:

kafkaParams.put("enable.auto.commit", false); // 关闭自动提交
// 移除auto.commit.interval.ms配置,因为手动提交不需要
kafkaParams.put("max.poll.interval.ms", "600000"); // 增大到10分钟,避免重平衡超时
kafkaParams.put("session.timeout.ms", "180000"); // 保持不变
kafkaParams.put("heartbeat.interval.ms", "60000"); // 改为session timeout的1/3

2. 添加手动偏移量提交逻辑

在stream.print()之后添加foreachRDD逻辑,确保处理完每个批次后再提交偏移量,避免重平衡时的冲突:

// 手动提交偏移量
stream.foreachRDD(new VoidFunction<JavaRDD<ConsumerRecord<String, String>>>() {
    @Override
    public void call(JavaRDD<ConsumerRecord<String, String>> rdd) throws Exception {
        // 获取当前批次的偏移量范围
        OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges();
        // 提交偏移量到Kafka
        ((CanCommitOffsets) stream.inputDStream()).commitAsync(offsetRanges);
    }
});

3. 确保Spark配置生效

你提交命令中已经添加了--conf spark.streaming.kafka.consumer.cache.enabled=false,这个配置非常重要,它会禁用Spark对Kafka消费者的缓存,确保每个任务的消费者都能参与重平衡,一定要保留。

4. 验证提交命令

保持你的提交命令不变,确保SPARK_KAFKA_VERSION=0.10和Kerberos相关配置正确:

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

为什么这些改动有效?

  • 手动提交偏移量让你完全控制偏移量提交时机,避免重平衡时自动提交的冲突
  • 调整后的心跳和会话超时参数符合Kafka的重平衡协调逻辑,减少协调器误判
  • 增大max.poll.interval.ms给重平衡足够的完成时间,避免消费者在重平衡期间被踢出组

内容的提问来源于stack exchange,提问作者venkat sam

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:08:40