The bellow mentioned Apache Flink test code doesn't compile in Scala when I want to use KeyedOneInputStreamOperatorTestHarness[K, IN, OUT] with Int or Long as the key (K) parameter. The same code can be compiled with generic parameters of type STRING, INTEGER or java.lang.Long:
import org.apache.flink.api.java.functions.KeySelector
import org.apache.flink.api.scala.typeutils.Types
import org.apache.flink.streaming.api.functions.KeyedProcessFunction
import org.apache.flink.streaming.api.operators.KeyedProcessOperator
import org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness
import org.apache.flink.streaming.api.scala._
import org.junit.Test
case class TaxiRide(rideId: Int)
class RideKeySelector extends KeySelector[TaxiRide, Int] {
override def getKey(in: TaxiRide): Int = in.rideId
}
class RideTests {
private def setupHarness(function: KeyedProcessFunction[Int, TaxiRide, Int]): KeyedOneInputStreamOperatorTestHarness[Int, TaxiRide, Int] = {
val operator: KeyedProcessOperator[Int, TaxiRide, Int] = new KeyedProcessOperator(function)
val testHarness: KeyedOneInputStreamOperatorTestHarness[Int, TaxiRide, Int] =
new KeyedOneInputStreamOperatorTestHarness(operator, new RideKeySelector(), Types.INT)
testHarness.setup()
testHarness.open()
testHarness
}
}
Build error:
type mismatch;
found : org.example.RideKeySelector
required: org.apache.flink.api.java.functions.KeySelector[org.example.TaxiRide,Int]
new KeyedOneInputStreamOperatorTestHarness(operator, new RideKeySelector(), Types.INT)
I created the project from the following Maven archetype:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-quickstart-scala</artifactId>
<version>1.14.2</version>
</dependency>
I added the following test dependencies to pom.xml in IntelliJ IDEA:
<dependencies>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-test-utils_2.11</artifactId>
<version>1.14.2</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-runtime</artifactId>
<version>1.14.2</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.11</artifactId>
<version>1.14.2</version>
<scope>test</scope>
<classifier>tests</classifier>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-scala_2.11</artifactId>
<version>1.14.2</version>
<scope>test</scope>
<classifier>tests</classifier>
</dependency>
</dependencies>
I tried to replace the key selector with a lambda:
val testHarness = new KeyedOneInputStreamOperatorTestHarness(operator, (ride: TaxiRide) => ride.rideId, Types.INT)
The build error changes to:
overloaded method constructor KeyedOneInputStreamOperatorTestHarness with alternatives:
(x$1: org.apache.flink.streaming.api.operators.OneInputStreamOperator[IN,OUT],x$2: org.apache.flink.api.java.functions.KeySelector[IN,K],x$3: org.apache.flink.api.common.typeinfo.TypeInformation[K],x$4: org.apache.flink.runtime.operators.testutils.MockEnvironment)org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness[K,IN,OUT] <and>
(x$1: org.apache.flink.streaming.api.operators.OneInputStreamOperator[IN,OUT],x$2: org.apache.flink.api.java.functions.KeySelector[IN,K],x$3: org.apache.flink.api.common.typeinfo.TypeInformation[K])org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness[K,IN,OUT] <and>
(x$1: org.apache.flink.streaming.api.operators.StreamOperatorFactory[OUT],x$2: org.apache.flink.api.java.functions.KeySelector[IN,K],x$3: org.apache.flink.api.common.typeinfo.TypeInformation[K])org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness[K,IN,OUT] <and>
(x$1: org.apache.flink.streaming.api.operators.StreamOperatorFactory[OUT],x$2: org.apache.flink.api.java.functions.KeySelector[IN,K],x$3: org.apache.flink.api.common.typeinfo.TypeInformation[K],x$4: Int,x$5: Int,x$6: Int)org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness[K,IN,OUT] <and>
(x$1: org.apache.flink.streaming.api.operators.OneInputStreamOperator[IN,OUT],x$2: org.apache.flink.api.java.functions.KeySelector[IN,K],x$3: org.apache.flink.api.common.typeinfo.TypeInformation[K],x$4: Int,x$5: Int,x$6: Int)org.apache.flink.streaming.util.KeyedOneInputStreamOperatorTestHarness[K,IN,OUT]
cannot be applied to (org.apache.flink.streaming.api.operators.KeyedProcessOperator[Int,org.example.TaxiRide,Int], org.example.TaxiRide => Integer, org.apache.flink.api.common.typeinfo.TypeInformation[Integer])
testHarness = new KeyedOneInputStreamOperatorTestHarness(operator, (ride: TaxiRide) => ride.rideId, Types.INT)
It looks like that it cannot accept Scala types, why is it so?