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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 09:31:34