class FilterFilter extends FilterFunction[String] { override def filter(value: String): Boolean = { value.contains("flink") } } val flinkTweets = tweets.filter(new FlinkFilter)还可以将函数实现成匿名类
val flinkTweets = tweets.filter( new RichFilterFunction[String] { override def filter(value: String): Boolean = { value.contains("flink") } } )我们 filter 的字符串"flink"还可以当作参数传进去。
class MyFilter(keyWord: String) extends FilterFunction[SensorReading]{ override def filter(value: SensorReading): Boolean = { value.id.contains(keyWord) } } val dataStream: DataStream[SensorReading] = inputStream .map( data => { val arr = data.split(",") SensorReading(arr(0), arr(1).toLong, arr(2).toDouble) } )
val tweets: DataStream[String] = ... val flinkTweets = tweets.filter(_.contains("flink"))
// 富函数,可以获取到运行时上下文,还有一些生命周期 class MyRichMap extends RichMapFunction[SensorReading, String]{ override def open(parameters: Configuration): Unit = { //做一些初始化操作。比如map方法需要交互数据库,数据库连接可以在open里边做 //getRuntimeContext() } override def map(value: SensorReading): String = { value.id + " temperature" } override def close(): Unit = { //map调用完之后。一般做收尾工作,比如关闭连接,或者清空状态 } }