如何在Scala中定义Flink Cassandra Sink的Mapper及相关配置
Flink Cassandra Sink Mapper配置的Scala等效实现
Java版本中.setMapperOptions(() -> new Mapper.Option[]{Mapper.Option.saveNullFields(true)})的等效Scala写法根据使用的Scala版本分为两种:
Scala 2.12+ 版本(支持SAM转换)
直接使用Lambda表达式即可,Scala会自动完成Java函数式接口的适配,Array类型直接对应Java数组类型,完整代码如下:
import com.datastax.driver.mapping.Mapper import org.apache.flink.streaming.connectors.cassandra.CassandraSink CassandraSink.addSink(result) .setHost("127.0.0.1") .setMapperOptions(() => Array[Mapper.Option](Mapper.Option.saveNullFields(true))) .build()
Scala 2.11 版本(无原生SAM转换支持)
需要显式实现方法要求的java.util.function.Supplier接口,写法如下:
import java.util.function.Supplier import com.datastax.driver.mapping.Mapper import org.apache.flink.streaming.connectors.cassandra.CassandraSink CassandraSink.addSink(result) .setHost("127.0.0.1") .setMapperOptions(new Supplier[Array[Mapper.Option]] { override def get(): Array[Mapper.Option] = Array(Mapper.Option.saveNullFields(true)) }) .build()
注意事项
- 确保引入的
Mapper类来自com.datastax.driver.mapping.Mapper包,不要和其他版本Cassandra驱动的同名类混淆 - Scala的
Array[T]在编译后直接对应Java的T[]数组类型,不需要额外做集合转换
内容的提问来源于stack exchange,提问作者Sjoerd222888
相关产品推荐
相关产品推荐

