同Consumer Group ID的两个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的处理逻辑没适配这个过程:
heartbeat.interval.ms设置过小(3000ms),和session.timeout.ms(180000ms)的比例严重偏离Kafka推荐的1/3原则,可能导致协调器误判消费者状态- 自动提交偏移量(
enable.auto.commit=true)在重平衡时容易出现偏移量提交冲突 - 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

