Scala 2.11结合Flink 1.17消费Kafka消息遇版本兼容问题求助
解决方案:Scala 2.11 + Flink流处理消费Kafka的依赖兼容问题
核心问题根源
Flink 1.15及以后版本(包括1.17)已停止对Scala 2.11的支持,官方仅维护Scala 2.12和2.13的适配包。这就是你引入flink-scala-streaming 1.17失败的原因;而混用1.14和1.17的依赖必然会引发跨版本冲突(API不兼容、类加载异常等)。
可行的解决路径
路径1:升级Scala版本到2.12/2.13(推荐)
这是最稳妥的方案,完全适配Flink 1.17的全部特性,包括新版Kafka Source。
- 调整项目Scala版本为2.12.x(Flink 1.17官方推荐),然后引入标准Flink Kafka依赖:
<!-- Maven示例 --> <dependencies> <!-- Flink核心API --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.0</version> <scope>provided</scope> </dependency> <!-- Flink Kafka连接器 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.0</version> </dependency> <!-- Scala API(如果需要用Scala编写) --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-scala</artifactId> <version>1.17.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-scala</artifactId> <version>1.17.0</version> <scope>provided</scope> </dependency> </dependencies>
路径2:固定使用Flink 1.14.x(兼容Scala 2.11)
如果无法升级Scala,只能放弃Flink 1.17,全量使用兼容Scala 2.11的最高Flink版本(1.14.x),避免跨版本混用:
- 统一所有Flink依赖版本为1.14.6(1.14系列最后一个稳定版):
<!-- Maven示例 --> <dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-scala_2.11</artifactId> <version>1.14.6</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_2.11</artifactId> <version>1.14.6</version> </dependency> </dependencies>
- 注意:Flink 1.14的Kafka Source使用的是旧版API(
FlinkKafkaConsumer),而非1.17的新Source API,代码写法需要对应调整:
// Flink 1.14 Scala消费Kafka示例 import org.apache.flink.streaming.api.scala._ import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer import java.util.Properties val env = StreamExecutionEnvironment.getExecutionEnvironment val props = new Properties() props.setProperty("bootstrap.servers", "localhost:9092") props.setProperty("group.id", "flink-consumer-group") val kafkaSource = new FlinkKafkaConsumer[String]("topic-name", new SimpleStringSchema(), props) val stream = env.addSource(kafkaSource) stream.print() env.execute("Kafka Streaming Job")
路径3:避免混用Scala API,直接使用Java API(不推荐,但可应急)
如果坚持用Flink 1.17 + Scala 2.11,可以放弃flink-scala-streaming,直接基于Flink Java API编写消费逻辑,Scala代码可以调用Java API:
- 只引入Flink Java和Kafka连接器依赖:
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>1.17.0</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>1.17.0</version> </dependency>
- 但这种方式会丢失Scala API的语法糖,且后续维护风险高,不建议长期使用。
内容的提问来源于stack exchange,提问作者Sachin Prajapati
相关产品推荐
相关产品推荐

