使用DataStax Connector读取Cassandra TIME类型遇编码器及转换异常
我正尝试使用DataStax Connector将Cassandra表读取到Spark中,该表有2列使用TIME数据类型。我在数据集里使用java.sql.Time作为对应类型,但Spark抛出以下异常:
Exception in thread "main" java.lang.UnsupportedOperationException: No Encoder found for java.sql.Time
- field (class: "java.sql.Time", name: "start")
- root class: "model.Trigger"
我已尝试通过Kryo注册Time类,但没有效果。我想知道是否应该使用其他类来对接Cassandra的TIME类型,或者是否存在作用域问题——我在主方法中注册Kryo类,但在另一个方法中从Cassandra获取数据(通过传递由配置生成的会话)。
感谢!
4月12日更新
我编写了一个自定义映射器来解析Time,但Spark抛出以下异常:
Exception in thread "main" java.lang.IllegalArgumentException: Unsupported type: java.sql.Time
at com.datastax.spark.connector.types.TypeConverter$.forCollectionType(TypeConverter.scala:1025)
at com.datastax.spark.connector.types.TypeConverter$.forType(TypeConverter.scala:1038)
at com.datastax.spark.connector.types.TypeConverter$.forType(TypeConverter.scala:1057)
对应的映射器代码如下:
object ColumnMappers { private object LongToTimeConverter extends TypeConverter[Time] { override def targetTypeTag: universe.TypeTag[Time] = typeTag[Time] override def convertPF: PartialFunction[Any, Time] = { case l: Long => Time.valueOf(LocalTime.ofNanoOfDay(l)) } } TypeConverter.registerConverter(LongToTimeConverter) private object TimeToLongConverter extends TypeConverter[Long] { override def targetTypeTag: universe.TypeTag[Long] = typeTag[Long] override def convertPF: PartialFunction[Any, Long] = { case t: Time => t.toLocalTime.toNanoOfDay } } TypeConverter.registerConverter(TimeToLongConverter) }
解决方案
方案1:改用java.time.LocalTime(推荐)
Spark对Java 8的时间API有原生Encoder支持,且DataStax Connector默认会将Cassandra的TIME类型(存储为当日纳秒数的Long值)直接映射到java.time.LocalTime,无需额外配置:
- 修改
model.Trigger类中对应字段的类型为java.time.LocalTime - 直接读取数据即可,无需自定义转换器或注册Kryo
方案2:保留java.sql.Time并修复自定义转换器
如果必须使用java.sql.Time,需要解决两个核心问题:
- Spark Encoder缺失:需要为
java.sql.Time提供自定义Encoder - DataStax转换器注册时机:确保转换器在Spark会话初始化前被加载
步骤1:添加java.sql.Time的Spark Encoder
在Spark代码中添加自定义Encoder:
import org.apache.spark.sql.Encoder import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder import java.sql.Time implicit val timeEncoder: Encoder[Time] = ExpressionEncoder()
步骤2:修复DataStax转换器的注册与使用
Scala单例对象是懒加载的,需确保ColumnMappers在SparkSession创建前被初始化(比如在main方法开头调用ColumnMappers)。或者改用ColumnMapper显式映射实体类:
import com.datastax.spark.connector.mapper.ColumnMapper import com.datastax.spark.connector.types.{TypeConverter, TimeType} import java.sql.Time import java.time.LocalTime class TriggerColumnMapper extends ColumnMapper[model.Trigger] { override def toCassandraType(clazz: Class[_]) = super.toCassandraType(clazz) match { case _ if clazz == classOf[Time] => TimeType case t => t } override def converterToCassandra(targetType: Class[_]) = super.converterToCassandra(targetType) match { case _ if targetType == classOf[Time] => TypeConverter.converterToCassandra[Long](classOf[Long]).compose[Time](_.toLocalTime.toNanoOfDay) case c => c } override def converterToScala(targetType: Class[_]) = super.converterToScala(targetType) match { case _ if targetType == classOf[Time] => TypeConverter.converterToScala[Long](classOf[Long]).andThen(l => Time.valueOf(LocalTime.ofNanoOfDay(l))) case c => c } }
然后在读取数据时指定该映射器:
val spark = SparkSession.builder() .appName("CassandraSparkDemo") .config("spark.cassandra.connection.host", "your-cassandra-host") .getOrCreate() import spark.implicits._ import com.datastax.spark.connector._ // 设置默认映射器 spark.sparkContext.setDefaultColumnMapper(new TriggerColumnMapper()) // 读取数据并转换为Dataset[model.Trigger] val triggerDS = spark.read .format("org.apache.spark.sql.cassandra") .options(Map("table" -> "trigger", "keyspace" -> "your-keyspace")) .load() .as[model.Trigger]
关于Kryo注册的误区
Kryo主要用于RDD的序列化,而Dataset依赖的是Encoder机制,所以注册Kryo无法解决No Encoder found的问题,这也是你之前尝试无效的原因。
内容的提问来源于stack exchange,提问作者maxwellray

