Scala调用Flink fromCollection报无法解析重载方法错误如何解决
报错根因
报错由两个问题共同导致:
- 导入类错误:当前代码引入的是Java版Flink API的
org.apache.flink.streaming.api.environment.StreamExecutionEnvironment,该类的fromCollection所有重载方法仅支持接收Java集合体系下的类型(实现java.util.Collection接口、或Java Iterator类型),但代码中传入的参数是Scala标准库的不可变List,不属于Java集合体系,编译器找不到匹配的重载方法。 - 泛型不匹配:传入集合的元素类型是自定义类
RiverWaterRegime,但代码将返回值声明为DataStreamSource[String],类型完全不匹配,也会导致编译校验失败。
修复方案
根据你使用的API类型二选一即可,你的项目已经引入了flink-streaming-scala依赖,优先选择Scala API方案。
方案1:使用Flink Scala API(Scala项目推荐)
Flink为Scala提供了专属API层,原生支持Scala集合类型,不需要手动做集合转换:
- 替换原有Java版执行环境的导入,改为引入Scala版的执行环境和相关类
- 导入Flink Scala API提供的隐式转换包,这是Scala API正常工作的必要依赖
- 修正返回值泛型,和集合元素类型保持一致
修正后的代码如下:
package com.mycode.learnflink.ts.sync import com.mycode.learnflink.model.datasourcesync.domain.RiverWaterRegime // 替换为Scala版API的类 import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment} // 导入Scala API必需的隐式转换 import org.apache.flink.streaming.api.scala._ class SyncWater { var cachedRiverWaterRegimes: List[RiverWaterRegime] = List() private val BATCH_SIZE = 1000 private var size = 0 def flinkProcess(list:List[RiverWaterRegime] ): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment env.setParallelism(6) // 泛型和集合元素类型一致,Scala版fromCollection原生支持Scala List val source: DataStream[RiverWaterRegime] = env.fromCollection(list) } }
方案2:保留Java版Flink API
如果需要继续使用Java版的StreamExecutionEnvironment,手动将Scala集合转换为Java集合即可:
- 导入Scala标准库提供的集合转换工具
- 调用
asJava方法将Scala List转为Java List后传入 - 修正返回值泛型和元素类型一致
修正后的代码如下:
package com.mycode.learnflink.ts.sync import com.mycode.learnflink.model.datasourcesync.domain.RiverWaterRegime import org.apache.flink.streaming.api.datastream.DataStreamSource import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment // 导入Scala到Java的集合转换工具 import scala.collection.JavaConverters._ class SyncWater { var cachedRiverWaterRegimes: List[RiverWaterRegime] = List() private val BATCH_SIZE = 1000 private var size = 0 def flinkProcess(list:List[RiverWaterRegime] ): Unit = { val env = StreamExecutionEnvironment.getExecutionEnvironment env.setParallelism(6) // 转换为Java集合后传入,修正泛型类型 val source:DataStreamSource[RiverWaterRegime] = env.fromCollection(list.asJava) } }
注:你当前pom.xml中的Scala版本、Flink版本匹配关系正常,不需要修改依赖配置。
内容的提问来源于stack exchange,提问作者Criwran
相关产品推荐
相关产品推荐

