【问题标题】:How to use java.time.LocalDate in Cassandra query from Spark?如何在来自 Spark 的 Cassandra 查询中使用 java.time.LocalDate?
【发布时间】:2017-02-22 14:43:54
【问题描述】:

我们在 Cassandra 中有一个表,其列 start_time 类型为 date

当我们执行以下代码时:

val resultRDD = inputRDD.joinWithCassandraTable(KEY_SPACE,TABLE)
   .where("start_time = ?", java.time.LocalDate.now)

我们得到以下错误:

com.datastax.spark.connector.types.TypeConversionException: Cannot convert object 2016-10-13 of type class java.time.LocalDate to com.datastax.driver.core.LocalDate.
at com.datastax.spark.connector.types.TypeConverter$$anonfun$convert$1.apply(TypeConverter.scala:45)
at com.datastax.spark.connector.types.TypeConverter$$anonfun$convert$1.apply(TypeConverter.scala:43)
at com.datastax.spark.connector.types.TypeConverter$LocalDateConverter$$anonfun$convertPF$14.applyOrElse(TypeConverter.scala:449)
at com.datastax.spark.connector.types.TypeConverter$class.convert(TypeConverter.scala:43)
at com.datastax.spark.connector.types.TypeConverter$LocalDateConverter$.com$datastax$spark$connector$types$NullableTypeConverter$$super$convert(TypeConverter.scala:439)
at com.datastax.spark.connector.types.NullableTypeConverter$class.convert(TypeConverter.scala:56)
at com.datastax.spark.connector.types.TypeConverter$LocalDateConverter$.convert(TypeConverter.scala:439)
at com.datastax.spark.connector.types.TypeConverter$OptionToNullConverter$$anonfun$convertPF$29.applyOrElse(TypeConverter.scala:788)
at com.datastax.spark.connector.types.TypeConverter$class.convert(TypeConverter.scala:43)
at com.datastax.spark.connector.types.TypeConverter$OptionToNullConverter.com$datastax$spark$connector$types$NullableTypeConverter$$super$convert(TypeConverter.scala:771)
at com.datastax.spark.connector.types.NullableTypeConverter$class.convert(TypeConverter.scala:56)
at com.datastax.spark.connector.types.TypeConverter$OptionToNullConverter.convert(TypeConverter.scala:771)
at com.datastax.spark.connector.writer.BoundStatementBuilder$$anonfun$8.apply(BoundStatementBuilder.scala:93)

我已经尝试根据documentation注册自定义转换器:

object JavaLocalDateToCassandraLocalDateConverter extends TypeConverter[com.datastax.driver.core.LocalDate] {
  def targetTypeTag = typeTag[com.datastax.driver.core.LocalDate]
  def convertPF = { 
      case ld: java.time.LocalDate => com.datastax.driver.core.LocalDate.fromYearMonthDay(ld.getYear, ld.getMonthValue, ld.getDayOfMonth) 
      case _ => com.datastax.driver.core.LocalDate.fromYearMonthDay(1971, 1, 1) 
  }
}

object CassandraLocalDateToJavaLocalDateConverter extends TypeConverter[java.time.LocalDate] {
  def targetTypeTag = typeTag[java.time.LocalDate]
  def convertPF = { case ld: com.datastax.driver.core.LocalDate => java.time.LocalDate.of(ld.getYear(), ld.getMonth(), ld.getDay()) 
                    case _ => java.time.LocalDate.now 
  }
}

TypeConverter.registerConverter(JavaLocalDateToCassandraLocalDateConverter)
TypeConverter.registerConverter(CassandraLocalDateToJavaLocalDateConverter)

但这并没有帮助。

如何在从 Spark 执行的 Cassandra 查询中使用 JDK8 日期/时间类?

【问题讨论】:

  • 内联转换后直接传递DataStax日期怎么样?你不喜欢吗?
  • 首先 - 在代码的其他部分,我们使用来自 JDK8 日期/时间的类,所以我不想每次都转换它。其次 - 即使我通过 DataStax LocalDate 我得到了object not serializable (class: com.datastax.driver.core.LocalDate, value: 2016-10-13)

标签: apache-spark cassandra apache-spark-sql spark-cassandra-connector


【解决方案1】:

我认为在这样的 where 子句中做的最简单的事情就是调用

sc
 .cassandraTable("test","test")
 .where("start_time = ?", java.time.LocalDate.now.toString)
 .collect`

只需传入字符串,因为这将是一个定义明确的转换。

TypeConverters 中似乎存在一个问题,您的转换器没有优先于内置转换器。我会快速浏览一下。

--编辑--

注册的转换器似乎没有正确地转移到执行者。在本地模式下,代码按预期工作,这让我认为这是一个序列化问题。我会为这个问题在 Spark Cassandra 连接器上开一张票。

【讨论】:

  • 感谢您的回复 - 它有效!在@shankar 回复之后,我想到了这个解决方案。但对我来说,我需要将java.time.LocalDate 转换为String,然后com.datastax.spark.connector.types.TypeConverter.LocalDateConverter 将其转换为com.datastax.driver.core.LocalDate,这看起来有点不太理想。
【解决方案2】:

Cassandra 日期格式为yyyy-MM-dd HH:mm:ss.SSS

所以你可以使用下面的代码,如果你使用 Java 8 将 Cassandra 日期转换为LocalDate,那么你可以做你的逻辑。

val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss.SSS")
val dateTime = LocalDateTime.parse(cassandraDateTime, formatter);

或者您可以将 LocalDate 转换为 Cassandra 日期格式并检查它。

【讨论】:

猜你喜欢
  • 1970-01-01
  • 2019-08-11
  • 2017-10-22
  • 2018-01-01
  • 2017-08-19
  • 2021-09-26
  • 2018-09-02
  • 1970-01-01
  • 2019-11-13
相关资源
最近更新 更多