Spark2.1+Kafka10偏移量管理报错:无法访问kafka010包的Assign
看起来你遇到的问题是Spark 2.1中Kafka 0.10模块的Assign类访问权限问题,其实核心原因是在Spark 2.1的kafka010 API里,Assign并不是可以直接实例化的公开类,而是ConsumerStrategies对象下的一个工厂方法。另外你的代码里还有几个小细节需要调整,下面一步步帮你解决:
问题根源
在Spark 2.1的org.apache.spark.streaming.kafka010包中,Assign是ConsumerStrategies对象的内部实现类,没有公开的构造方法,直接尝试Assign[String, String](...)会因为访问权限问题报错。正确的做法是使用ConsumerStrategies.assign()方法来创建消费策略。
解决方案步骤
1. 修正消费策略的创建方式
把直接实例化Assign的代码,替换为ConsumerStrategies.assign()调用,同时确保你已经正确导入ConsumerStrategies。
2. 修正fromOffsets的类型定义
你的fromOffsets定义为Map[Object, Long],这会导致类型不匹配,应该明确为Map[TopicPartition, Long],这样和Kafka API更兼容。
3. 确认依赖版本匹配
确保你的spark-streaming-kafka-0-10依赖版本和Spark 2.1完全一致,比如如果用Maven,依赖应该是:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.11</artifactId> <version>2.1.0</version> </dependency>
修正后的完整代码
import com.typesafe.config.{ConfigFactory, ConfigValueFactory} import org.apache.kafka.clients.consumer.ConsumerRecord import org.apache.kafka.common.TopicPartition import org.apache.spark.sql.{Row, SparkSession} import org.apache.spark.streaming.dstream.InputDStream import org.apache.spark.streaming.kafka010._ import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.sql.functions._ import org.apache.spark.sql.types.{StringType, StructField, StructType} object offTest { def main(args: Array[String]) { val Array(impalaHost, brokers, topics, consumerGroupId, ssl, truststoreLocation, truststorePassword, wInterval) = args val sparkSession = SparkSession.builder .config("spark.hadoop.parquet.enable.summary-metadata", "false") .enableHiveSupport() .getOrCreate val ssc = new StreamingContext(sparkSession.sparkContext, Seconds(wInterval.toInt)) val isUsingSsl = ssl.toBoolean // Create direct kafka stream with brokers and topics val topicsSet = topics.split(",").toSet val commonParams = Map[String, Object]( "bootstrap.servers" -> brokers, "security.protocol" -> (if (isUsingSsl) "SASL_SSL" else "SASL_PLAINTEXT"), "sasl.kerberos.service.name" -> "kafka", "auto.offset.reset" -> "latest", "key.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "value.deserializer" -> "org.apache.kafka.common.serialization.StringDeserializer", "group.id" -> consumerGroupId, "enable.auto.commit" -> (false: java.lang.Boolean) ) val additionalSslParams = if (isUsingSsl) { Map( "ssl.truststore.location" -> truststoreLocation, "ssl.truststore.password" -> truststorePassword ) } else { Map.empty } val kafkaParams = commonParams ++ additionalSslParams // 修正fromOffsets的类型为Map[TopicPartition, Long] val fromOffsets: Map[TopicPartition, Long] = Map( new TopicPartition(topics, 4) -> 4807048129L ) // 使用ConsumerStrategies.assign替代直接实例化Assign val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.assign[String, String](fromOffsets.keys.toList, kafkaParams, fromOffsets) ) stream.foreachRDD { rdd => val offsetRanges = rdd.asInstanceOf[HasOffsetRanges].offsetRanges rdd.foreachPartition { iter => val o: OffsetRange = offsetRanges(TaskContext.get.partitionId) println(s"${o.topic} ${o.partition} ${o.fromOffset} ${o.untilOffset}") // I will insert those values to other database later } } val data= stream.map(record => (record.key, record.value)) data.foreachRDD(rdd1 => { val value = rdd1.map(x => x._2) if (!value.isEmpty()) { value.foreach(println) } else {println("no data")} }) ssc.start() ssc.awaitTermination() } }
关键修改说明
- 替换
Assign[String, String](...)为ConsumerStrategies.assign[String, String](...):这是Spark 2.1 Kafka 0.10 API的标准用法,ConsumerStrategies提供了创建消费策略的公开方法。 - 明确
fromOffsets的类型:避免Object类型带来的隐式转换问题,确保类型安全。 - 确保导入了
org.apache.spark.streaming.kafka010._:这个导入包含了ConsumerStrategies对象,所以不需要额外单独导入。
如果还是遇到包冲突问题,可以尝试清理项目依赖(比如Maven的clean install或者SBT的clean compile),排除掉冲突的Kafka或Spark依赖,确保只有Spark 2.1对应的kafka010依赖被引入。
内容的提问来源于stack exchange,提问作者Humayra Raisha

