
import org.apache.flink.api.scala.createTypeInformation
import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment
object filterTest {
def main(args: Array[String]): Unit = {
//create env
val env = StreamExecutionEnvironment.getExecutionEnvironment
//create ds
val ds = env.fromElements("hadoop is good", "flink is faster", "flink is better")
//transformation
val filteredDs = ds.filter(_.contains("flink"))
filteredDs.print()
env.execute()
}
}
