Spark写入Hive TimestampType时纳秒被截断的解决方案咨询
解决Spark写入Hive时Timestamp纳秒截断的问题
这个问题我之前帮别人排查过,本质是Spark和Hive在Timestamp精度支持上的兼容性问题,结合你提到的ORC和CSV两种格式,给你几个可行的解决思路:
一、针对ORC格式:升级Hive版本+开启纳秒支持
Hive对ORC格式的Timestamp精度支持是分版本的:
- Hive 2.2及之前的版本,ORC里的Timestamp最多只支持微秒(6位小数),所以你的纳秒后3位会被截断;
- 从Hive 3.0开始,ORC 0.12格式正式支持纳秒级Timestamp。
如果你的集群能升级到Hive 3.x+,只需要在Spark作业里添加以下配置即可:
val sparkConf = new SparkConf().setAppName("Test") .set("hive.exec.orc.default.format.version", "0.12") // 指定ORC版本支持纳秒 .set("spark.sql.orc.impl", "native") // 使用原生ORC实现 .set("spark.sql.orc.enableVectorizedReader", "true") // 开启向量化读取,提升性能
这样写入ORC格式的Hive表时,就能完整保留纳秒部分了。
二、针对CSV格式:自定义时间输出格式
Spark写CSV时,默认的timestampFormat只输出到毫秒(3位小数),所以你看到的是2018-03-20T13:04:20.123Z。只需要在写入时指定包含纳秒的格式即可:
df.write .mode("overwrite") .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSSSSSSSS'Z'") // 9位小数对应纳秒 .csv("/path/to/your/csv/output")
配置后CSV里的Timestamp就会显示为2018-03-20T13:04:20.123456789Z,完整保留纳秒信息。
三、临时替代方案:将Timestamp转字符串存储
如果暂时无法升级Hive版本,最稳妥的临时方案是把Timestamp转换为字符串类型存储,读取时再转回去:
写入时转换
import org.apache.spark.sql.functions._ df.withColumn("timestamp_full", date_format(col("dt"), "yyyy-MM-dd HH:mm:ss.SSSSSSSSS")) .drop("dt") // 替换原Timestamp字段,或者保留原字段 .write .saveAsTable("your_hive_table")
读取时转换
sqlContext.read.table("your_hive_table") .withColumn("dt", to_timestamp(col("timestamp_full"), "yyyy-MM-dd HH:mm:ss.SSSSSSSSS"))
这种方式完全不会丢失精度,缺点是无法直接用Timestamp的函数做查询,需要先转换。
给你的测试代码补充写入逻辑
结合你的测试代码,这里给你补全写入Hive的完整示例,方便你验证:
import org.apache.spark.SparkConf import org.apache.spark.SparkContext import org.apache.spark.sql.types._ import org.apache.spark.sql.Row import java.math.BigDecimal import org.apache.spark.sql.functions._ object testDateAndDecimal { def main(args: Array[String]): Unit = { execute; } private def execute: Unit = { val sparkConf = new SparkConf().setAppName("Test") // ORC纳秒支持配置 .set("hive.exec.orc.default.format.version", "0.12") .set("spark.sql.orc.impl", "native") val sc = new SparkContext(sparkConf) val sqlContext = new org.apache.spark.sql.SQLContext(sc) // Define DataTypes val datetimestring: String = "2018-03-20 13:04:20.123456789" val dt = java.sql.Timestamp.valueOf(datetimestring) val price = new BigDecimal("1234567890.1234567899") // 构造测试DataFrame val schema = StructType(Seq( StructField("id", IntegerType), StructField("dt", TimestampType), StructField("price", DecimalType(18, 9)) )) val data = Seq(Row(1, dt, price)) val df = sqlContext.createDataFrame(sc.parallelize(data), schema) // 写入ORC格式Hive表 df.write .mode("overwrite") .format("orc") .saveAsTable("test_timestamp_orc") // 写入CSV格式(保留纳秒) df.write .mode("overwrite") .option("timestampFormat", "yyyy-MM-dd'T'HH:mm:ss.SSSSSSSSS'Z'") .csv("/tmp/test_timestamp_csv") sc.stop() } }
内容的提问来源于stack exchange,提问作者Paras Gupta
相关产品推荐
相关产品推荐

