Spark 2.1+Kafka 2.1环境下Spark Streaming替代KafkaUtils的对接方案咨询
问题解答
参考内容排查说明
你没有遗漏相关参考内容:你之前使用的KafkaUtils.createStream是Spark 1.x配套的spark-streaming-kafka-0-8集成包专属API,该集成包从Spark 2.0版本开始就已停止官方维护,针对Kafka 0.10及以上版本(包括你升级后的Kafka 2.1),官方统一提供了新的spark-streaming-kafka-0-10集成方案,API接口和旧版完全不同。
适配Spark 2.1 + Kafka 2.1的最优实现方案(Java语言)
推荐使用官方原生的Direct Stream对接方案,相比旧版Receiver模式资源利用率更高、数据一致性更强,不需要依赖第三方组件,具体实现步骤如下:
- 第一步:引入匹配版本的依赖
以Maven为例,Cloudera 6.2的Spark 2.1默认基于Scala 2.11编译,依赖配置参考如下:<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-streaming-kafka-0-10_2.11</artifactId> <!-- CDH环境使用对应CDH版本号,开源环境直接写2.1.0即可 --> <version>2.1.0-cloudera2</version> </dependency> - 第二步:核心业务代码实现
新API直接兼容Kafka 2.1的所有配置参数,原有业务处理逻辑可以完全复用,示例代码如下:import org.apache.spark.SparkConf; import org.apache.spark.streaming.Durations; import org.apache.spark.streaming.api.java.JavaStreamingContext; import org.apache.spark.streaming.kafka010.ConsumerStrategies; import org.apache.spark.streaming.kafka010.KafkaUtils; import org.apache.spark.streaming.kafka010.LocationStrategies; import org.apache.kafka.common.serialization.StringDeserializer; import java.util.Collections; import java.util.HashMap; import java.util.Map; public class KafkaStreamTask { public static void main(String[] args) throws InterruptedException { // 初始化Spark Streaming上下文,批次间隔按需调整 SparkConf conf = new SparkConf().setAppName("KafkaStreamTask"); JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(5)); // Kafka消费者配置,所有原生Kafka消费者参数都可直接配置 Map<String, Object> kafkaParams = new HashMap<>(); kafkaParams.put("bootstrap.servers", "broker1:9092,broker2:9092"); kafkaParams.put("key.deserializer", StringDeserializer.class); kafkaParams.put("value.deserializer", StringDeserializer.class); kafkaParams.put("group.id", "your-consumer-group-id"); kafkaParams.put("auto.offset.reset", "latest"); // 关闭自动提交,处理完成后手动提交保证数据不丢不重 kafkaParams.put("enable.auto.commit", false); // 订阅指定Topic String targetTopic = "your-business-topic"; // 创建Direct流 var kafkaStream = KafkaUtils.createDirectStream( jssc, LocationStrategies.PreferConsistent(), ConsumerStrategies.Subscribe(Collections.singleton(targetTopic), kafkaParams) ); // 原有业务处理逻辑直接复用即可 kafkaStream.foreachRDD(rdd -> { rdd.foreach(record -> { String dataKey = record.key(); String dataValue = record.value(); // 此处插入你原有数据处理逻辑 }); // 业务处理完成后手动提交offset // OffsetRange[] offsetRanges = ((HasOffsetRanges) rdd.rdd()).offsetRanges(); // ((CanCommitOffsets) kafkaStream.inputDStream()).commitAsync(offsetRanges); }); jssc.start(); jssc.awaitTermination(); } } - 第三步:迁移注意事项
- 新的Direct模式不需要配置Receiver相关参数,资源占用比旧版低30%以上
- Kafka 2.1完全兼容0-10版本的消费者API,不需要额外做版本适配
- 如果集群开启Kerberos认证,直接在kafkaParams中增加安全协议、JAAS配置等原生Kafka消费者参数即可
- 如果你需要兼容旧版消费者组的offset,可以先将旧offset同步到新版消费者组,或者调整
auto.offset.reset参数为earliest从头消费
内容的提问来源于stack exchange,提问作者elvaras
相关产品推荐
相关产品推荐

