Spark与Kafka集成报错:无法解析KafkaUtils的createStream方法
问题根源
你用的是Spark Streaming Kafka 0.10版本的KafkaUtils,但createStream是旧版(0.8/0.9)API的方法,0.10版本已经废弃该方法,改用createDirectStream且参数格式完全不同。另外pom中spark-streaming-kafka-0-10依赖设为provided,本地编译时会找不到相关类,这也是报错的原因之一。
解决步骤
1. 修正Maven依赖
把spark-streaming-kafka-0-10_2.13的<scope>从provided改成compile,本地开发需要这个依赖完成编译:
<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.13</artifactId> <version>3.3.1</version> <!-- 移除provided,本地编译需引入该依赖 --> </dependency>
2. 重写Kafka流初始化代码
替换原createStream调用,改用0.10版本的createDirectStream,并配置正确参数:
- 先添加必要的导入类:
import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.kafka.clients.consumer.ConsumerRecord; import java.util.Collections;
- 再修改main方法中的流初始化逻辑:
// 配置Kafka消费者核心参数 Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", BOOTSTRAP_SERVER); kafkaParams.put("group.id", GROUP_ID); kafkaParams.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); kafkaParams.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // 指定要订阅的主题 Collection<String> topics = Collections.singletonList("Demo1"); // 创建DirectStream JavaInputDStream<ConsumerRecord<String, String>> stream = KafkaUtils.createDirectStream( streamingContext, LocationStrategies.PreferConsistent(), // 位置策略:优先将分区分配给最近的Executor ConsumerStrategies.Subscribe(topics, kafkaParams) // 订阅策略 ); // 从ConsumerRecord中提取key和value,转换为JavaPairDStream JavaPairDStream<String, String> messages = stream.mapToPair(record -> new Tuple2<>(record.key(), record.value()) );
3. 完整修改后的代码
import java.util.HashMap; import java.util.Map; import java.util.Collections; import org.apache.spark.SparkConf; import org.apache.spark.streaming.Duration; import org.apache.spark.streaming.api.java.JavaPairDStream; import org.apache.spark.streaming.api.java.JavaStreamingContext; import org.apache.spark.streaming.kafka010.KafkaUtils; import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.kafka.clients.consumer.ConsumerRecord; import scala.Tuple2; public class NumbersConsumer { public static final String BOOTSTRAP_SERVER = "localhost:9092"; public static final String GROUP_ID = "spark-group"; public static final int BATCH_DURATION = 5000; public static void main(String[] args) throws Exception { // 本地测试添加setMaster("local[*]"),生产环境部署时移除该配置 SparkConf sparkConf = new SparkConf().setAppName("demo").setMaster("local[*]"); JavaStreamingContext streamingContext = new JavaStreamingContext(sparkConf, new Duration(BATCH_DURATION)); // Kafka消费者参数配置 Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", BOOTSTRAP_SERVER); kafkaParams.put("group.id", GROUP_ID); kafkaParams.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); kafkaParams.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); // 订阅主题 Collection<String> topics = Collections.singletonList("Demo1"); // 创建DirectStream JavaInputDStream<ConsumerRecord<String, String>> stream = KafkaUtils.createDirectStream( streamingContext, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(topics, kafkaParams) ); // 提取消息的key和value JavaPairDStream<String, String> messages = stream.mapToPair(record -> new Tuple2<>(record.key(), record.value()) ); JavaPairDStream<String, String> percentage = messages.mapToPair( num -> new Tuple2<>("Percentage: ", num._2 + "%") ); percentage.print(); streamingContext.start(); streamingContext.awaitTermination(); } }
补充说明
- 本地测试时必须给
SparkConf添加setMaster("local[*]"),否则会因找不到集群资源报错;生产环境部署时需移除该配置,由集群管理器统一分配资源。 - Kafka 0.10版本的API基于Kafka原生消费者实现,比旧版
createStream更高效,也更贴合Kafka原生功能。 - 确保Kafka集群版本与依赖的
kafka-clients版本一致(当前均为3.3.1,无需调整)。
内容的提问来源于stack exchange,提问作者user20148665
相关产品推荐
相关产品推荐

