Apache Flink Scala中Join操作报错:.apply()方法不存在?
Flink 1.16窗口Join代码问题及解决
我参照Apache Flink 1.16版本官方文档编写窗口Join代码,官方示例代码如下:
stream.join(otherStream) .where(<KeySelector>) .equalTo(<KeySelector>) .window(<WindowAssigner>) .apply(<JoinFunction>);
我的Scala代码实现如下:
import org.apache.flink.streaming.api.windowing.assigners.SlidingEventTimeWindows private val result = streamA.join(streamB) .where(new MySelector().getKey) .equalTo(new MySelector().getKey) .window(SlidingEventTimeWindows.of(Time.minutes(60), Time.minutes(60)) .apply() // <- 找不到该方法!
在IntelliJ中没有.apply()的提示选项,排查后确认不是Flink本身的问题,而是window方法调用时缺少了闭合括号。
以下是我的build.sbt配置:
ThisBuild / version := "0.1.0-SNAPSHOT" ThisBuild / scalaVersion := "2.12.12" lazy val root = (project in file(".")) .settings( name := "apache-flink-example", libraryDependencies += "org.apache.flink" % "flink-streaming-scala_2.12" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-connector-kafka" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-connectors" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-core" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-java8" % "0.10.1", libraryDependencies += "org.apache.flink" % "flink-quickstart" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-libraries" % "1.16.1", libraryDependencies += "org.apache.flink" % "flink-clients" % "1.16.1", libraryDependencies += "org.json4s" %% "json4s-native" % "4.0.6", )
内容的提问来源于stack exchange,提问作者kev
相关产品推荐
相关产品推荐

