Spark Streaming 2.0.0集成Embedded Kafka启动失败求助
Spark Streaming 2.0.0 + scalatest-embedded-kafka启动失败?看这几个解决方案
你遇到的情况很典型——Spark 2.2.0能正常跑,换成2.0.0就出问题,核心原因肯定是版本兼容性不匹配,毕竟scalatest-embedded-kafka和Spark、Kafka的版本绑定得相当紧密。下面给你几个针对性的解决思路:
1. 精准匹配适配Spark 2.0.0的embedded-kafka版本
别盲目试版本了,直接锁定对应生态的版本分支:
- 针对Scala 2.11 + Spark 2.0.0 + Kafka 0.10.x的组合,直接用
scalatest-embedded-kafka_2.11的0.10.x系列版本(比如0.10.0.0),这个分支就是专门为Kafka 0.10打造的,和Spark 2.0.x的依赖不会出现冲突。 - 千万别碰1.x及以上的embedded-kafka版本,那些都是为Spark 2.3+适配的,和Spark 2.0.0的内部API差异极大,必然会触发启动异常。
2. 手动排查并解决依赖冲突
Spark 2.0.0自带的Kafka客户端版本可能和embedded-kafka引入的版本不一致,导致类加载冲突。可以这么处理:
- 用Maven的
mvn dependency:tree或者SBT的sbt dependencyTree命令,查看输出里kafka-clients的版本是否统一为0.10.x。 - 如果发现多版本冲突,就在embedded-kafka的依赖中排除自带的Kafka相关组件,再手动引入和Spark 2.0.0兼容的版本:
<dependency> <groupId>net.manub</groupId> <artifactId>scalatest-embedded-kafka_2.11</artifactId> <version>0.10.0.0</version> <scope>test</scope> <exclusions> <exclusion> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> </exclusion> <exclusion> <groupId>org.apache.kafka</groupId> <artifactId>kafka_2.11</artifactId> </exclusion> </exclusions> </dependency> <!-- 显式引入兼容的Kafka客户端 --> <dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>0.10.2.2</version> <scope>test</scope> </dependency>
3. 调整Spark Streaming的初始化代码适配2.0.0
Spark 2.0.0的Streaming API和2.2.0有细节差异,别直接照搬2.2.0的代码,要适配2.0.0的kafka010 API:
import org.apache.spark.streaming.kafka010._ import org.apache.spark.streaming.{Seconds, StreamingContext} // 初始化StreamingContext val ssc = new StreamingContext(sparkConf, Seconds(5)) // Kafka参数要对应0.10的格式 val kafkaParams = Map[String, Object]( "bootstrap.servers" -> "localhost:9092", "key.deserializer" -> classOf[StringDeserializer], "value.deserializer" -> classOf[StringDeserializer], "group.id" -> "test-group", "auto.offset.reset" -> "latest", "enable.auto.commit" -> (false: java.lang.Boolean) ) val topics = Array("test-topic") // 用Spark 2.0.0支持的DirectStream创建方式 val stream = KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) )
注意:这里的LocationStrategies和ConsumerStrategies都是kafka010包下的,别误用旧版本的KafkaUtils API。
4. 显式指定embedded-kafka的启动配置
有时候默认配置会出现端口冲突或者参数不兼容的问题,你可以在测试里手动指定端口和配置:
import net.manub.embeddedkafka.{EmbeddedKafka, EmbeddedKafkaConfig} // 自定义Kafka和ZK端口,避免冲突 val customConfig = EmbeddedKafkaConfig(kafkaPort = 9092, zooKeeperPort = 2181) // 用自定义配置启动EmbeddedKafka EmbeddedKafka.start()(customConfig) // 这里写你的测试逻辑 // ... // 测试结束后停止服务 EmbeddedKafka.stop()
内容的提问来源于stack exchange,提问作者Guille
相关产品推荐
相关产品推荐

