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

使用DataStax Connector读取Cassandra TIME类型遇编码器及转换异常

使用DataStax Connector读取Cassandra TIME类型字段到Spark的问题

我正尝试使用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,需要解决两个核心问题:

  1. Spark Encoder缺失:需要为java.sql.Time提供自定义Encoder
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 05:24:58