You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Scala中定义Flink Cassandra Sink的Mapper及相关配置

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 18:12:57