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

Spark2.1+Kafka10偏移量管理报错:无法访问kafka010包的Assign

解决Spark Streaming Kafka 0.10中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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:12:46