Spark写入Cassandra:JSON Timestamp转TimeUUID的报错及解决方法
解决Spark作业中JSON时间字符串转Cassandra TimeUUID的问题
看起来你踩了一个很常见的坑:JSON里的timestamp字段是ISO 8601格式的字符串(比如"2019-05-09T09:00:52.553+0000"),但你直接尝试把它转成Long或者UUID,自然会报NumberFormatException——因为这些方法只认数字格式的输入,不认时间字符串。
下面是一步步的解决方案,帮你正确完成转换并插入Cassandra:
1. 先解析ISO时间字符串为毫秒时间戳
首先需要把时间字符串转换成Long类型的毫秒数,这是生成TimeUUID的前提。你可以用Java 8的java.time库(推荐,Spark 2.3+支持很好)或者Joda-Time(兼容旧版Spark)来解析:
// 用Java 8 DateTime API解析时间字符串 import java.time.{Instant, LocalDate, ZoneOffset} import java.time.format.DateTimeFormatter def parseTimestampToMillis(timestampStr: String): Option[Long] = { try { // 匹配你的时间格式:yyyy-MM-dd'T'HH:mm:ss.SSSZ val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSZ") val instant = Instant.from(formatter.parse(timestampStr)) Some(instant.toEpochMilli()) } catch { case e: Exception => // 打印错误日志,跳过无效数据 println(s"Failed to parse timestamp: $timestampStr, error: ${e.getMessage}") None } } // 如果你用的是旧版Spark(依赖Joda-Time),可以用这个版本 // import org.joda.time.{DateTime, DateTimeZone} // import org.joda.time.format.DateTimeFormat // def parseTimestampToMillis(timestampStr: String): Option[Long] = { // try { // val formatter = DateTimeFormat.forPattern("yyyy-MM-dd'T'HH:mm:ss.SSSZ") // val dateTime = DateTime.parse(timestampStr, formatter).withZone(DateTimeZone.UTC) // Some(dateTime.getMillis) // } catch { // case e: Exception => // println(s"Failed to parse timestamp: $timestampStr, error: ${e.getMessage}") // None // } // }
2. 用Cassandra工具类生成TimeUUID
Cassandra的com.datastax.driver.core.utils.UUIDs.startOf()方法可以接收毫秒时间戳,生成对应的TimeUUID,这完全符合你的需求——因为TimeUUID本身就是包含时间信息的UUID。
3. 修改Spark RDD转换逻辑
把上面的解析方法集成到你的RDD处理流程里,同时优化JSON解析的代码(原来的json4s用法有点问题):
import com.datastax.driver.core.utils.UUIDs import org.json4s._ import org.json4s.JsonMethods._ // 你的case class(注意字段名和类型要和Cassandra表对应) case class CassandraFormat( eventID: String, sessionID: String, timeuuid: UUID, userID: String, event_date: LocalDate, // 如果用Cassandra的LocalDate,改成com.datastax.driver.core.LocalDate fullJson: String ) val allJson = rdd // 解析Kafka消息为JSON对象,同时保留原始JSON字符串 .map(rawJson => { implicit val formats = DefaultFormats (parse(rawJson).extract[Map[String, Any]], rawJson) }) // 过滤掉没有header的消息 .filter(_._1.contains("header")) // 提取header部分 .map { case (jsonMap, rawJson) => (jsonMap("header").asInstanceOf[Map[String, Any]], rawJson) } // 过滤掉缺少必要字段的消息 .filter(header => header.contains("userID") && header.contains("eventID") && header.contains("sessionID") && header.contains("timestamp") ) // 转换为CassandraFormat对象,同时处理时间转换 .flatMap { case (header, rawJson) => val timestampStr = header("timestamp").toString parseTimestampToMillis(timestampStr).map { timestampMillis => // 生成TimeUUID val timeuuid = UUIDs.startOf(timestampMillis) // 生成event_date(这里用Java 8 LocalDate,对应Cassandra的date类型) val eventDate = LocalDate.ofInstant(Instant.ofEpochMilli(timestampMillis), ZoneOffset.UTC) // 如果case class用的是Cassandra的LocalDate,替换成: // val eventDate = com.datastax.driver.core.LocalDate.fromMillisSinceEpoch(timestampMillis) CassandraFormat( eventID = header("eventID").toString, sessionID = header("sessionID").toString, timeuuid = timeuuid, userID = header("userID").toString, event_date = eventDate, fullJson = rawJson ) } }
关键注意事项
- 时区处理:解析时间时一定要用UTC时区,避免因为服务器时区差异导致时间偏移,保证存入Cassandra的时间和原始事件时间一致。
- 异常处理:用
Option和flatMap跳过解析失败的无效数据,避免单个坏消息导致整个Spark作业崩溃。 - 类型匹配:确保case class的字段类型和Cassandra表完全对应,比如
event_date如果是Cassandra的date类型,用com.datastax.driver.core.LocalDate或者Java 8的LocalDate都可以(Spark Cassandra Connector会自动转换)。
这样处理后,你就能正确把JSON里的时间字符串转成Cassandra的TimeUUID,顺利完成数据插入了。
内容的提问来源于stack exchange,提问作者dejanmarich
相关产品推荐
相关产品推荐

