Skip to content

Commit

Permalink
fixup
Browse files Browse the repository at this point in the history
  • Loading branch information
Clarkkkkk committed Dec 10, 2018
1 parent 6438c8d commit 43915c1
Showing 1 changed file with 3 additions and 2 deletions.
Expand Up @@ -23,7 +23,7 @@ import java.lang
import org.apache.flink.api.common.functions._
import org.apache.flink.api.java.io.ParallelIteratorInputFormat
import org.apache.flink.api.java.typeutils.TypeExtractor
import org.apache.flink.streaming.api.collector.selector.OutputSelector
import org.apache.flink.streaming.api.collector.selector.{OutputSelector, OutputSelectorWrapper}
import org.apache.flink.streaming.api.functions.{KeyedProcessFunction, ProcessFunction}
import org.apache.flink.streaming.api.functions.co.CoMapFunction
import org.apache.flink.streaming.api.graph.{StreamEdge, StreamGraph}
Expand Down Expand Up @@ -567,7 +567,8 @@ class DataStreamTest extends AbstractTestBase {
split.print()
val outputSelectors = env.getStreamGraph.getStreamNode(unionFilter.getId).getOutputSelectors
assert(1 == outputSelectors.size)
assert(outputSelector == outputSelectors.get(0))
assert(outputSelector ==
outputSelectors.get(0).asInstanceOf[OutputSelectorWrapper[Int]].getCurrentOutputSelector)

unionFilter.split(x => List("a")).print()
val moreOutputSelectors = env.getStreamGraph.getStreamNode(unionFilter.getId).getOutputSelectors
Expand Down

0 comments on commit 43915c1

Please sign in to comment.