Spark3.0.1中日期时间列多格式转换报错的解决方法
问题描述
我正在使用Spark3.0.1,现有如下CSV格式数据:
348702330256514,37495066290,9084849,33946,614677375609919,11-02-2018 0:00:00,GENUINE
348702330256514,37495066290,330148,33946,614677375609919,11-02-2018 00:00:00,GENUINE
...(其余数据行略)
可以看到倒数第二列的transaction_dt为时间戳数据,小时部分存在一位数和两位数两种格式(示例数据仅展示时间全零情况,实际包含其他时间)。
我尝试先将该列读取为String类型,再通过自定义方法转换为Timestamp类型,代码如下:
val schema = StructType( List( StructField("_corrupt_record", StringType) , StructField("card_id", LongType) , StructField("member_id", LongType) , StructField("amount", IntegerType) , StructField("postcode", IntegerType) , StructField("pos_id", LongType) , StructField("transaction_dt", StringType) , StructField("status", StringType) ) ) // format the timestamp column def format_time_column(timeStampCol: Column , formats: Seq[String] = Seq( "dd-MM-yyyy HH:mm:ss", "dd-MM-yyyy H:mm:ss" , "dd-MM-yyyy HH:m:ss", "dd-MM-yyyy H:m:ss")) ={ coalesce( formats.map(f => to_timestamp(timeStampCol, f)):_* ) } val cardTransaction = spark.read .format("csv") .option("header", false) .schema(schema) .option("path", cardTransactionFilePath) .option("columnNameOfCorruptRecord", "_corrupt_record") .load .withColumn("transaction_dt", format_time_column(col("transaction_dt"))) cardTransaction.cache() cardTransaction.show(5)
但运行代码后出现报错,问题核心:
- 含一位数小时的记录触发报错
- 仅列表中第一个格式生效,其余格式未被识别
to_timestamp遇到错误格式时抛出异常,而非返回null,导致coalesce无法正常工作
解决方案
1. 开启宽松时间解析模式
Spark 3.0+默认采用严格时间解析策略,不匹配格式直接抛出异常。需通过配置spark.sql.legacy.timeParserPolicy为LEGACY,让to_timestamp解析失败时返回null,而非报错,这样coalesce才能正常依次尝试不同格式。
在初始化SparkSession时添加配置:
val spark = SparkSession.builder() .appName("CardTransactionProcessing") .config("spark.sql.legacy.timeParserPolicy", "LEGACY") .getOrCreate()
2. 简化时间格式定义
实际上,H格式符本身支持一位数和两位数的小时,m同样支持一位数和两位数的分钟,无需定义多个重复格式。仅保留dd-MM-yyyy H:mm:ss即可匹配所有场景。
简化后的转换方法:
def format_time_column(timeStampCol: Column, format: String = "dd-MM-yyyy H:mm:ss") = { to_timestamp(timeStampCol, format) }
3. 完整修改后的代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types._ import org.apache.spark.sql.functions._ object CardTransactionProcessor { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("CardTransactionProcessing") .config("spark.sql.legacy.timeParserPolicy", "LEGACY") .getOrCreate() import spark.implicits._ val schema = StructType( List( StructField("_corrupt_record", StringType), StructField("card_id", LongType), StructField("member_id", LongType), StructField("amount", IntegerType), StructField("postcode", IntegerType), StructField("pos_id", LongType), StructField("transaction_dt", StringType), StructField("status", StringType) ) ) // 简化后的时间格式转换方法 def format_time_column(timeStampCol: Column, format: String = "dd-MM-yyyy H:mm:ss") = { to_timestamp(timeStampCol, format) } val cardTransactionFilePath = "your_file_path_here" // 替换为实际文件路径 val cardTransaction = spark.read .format("csv") .option("header", false) .schema(schema) .option("path", cardTransactionFilePath) .option("columnNameOfCorruptRecord", "_corrupt_record") .load() .withColumn("transaction_dt", format_time_column(col("transaction_dt"))) cardTransaction.cache() cardTransaction.show(5) } }
额外说明
如果业务中确实存在多种完全不同的时间格式,开启LEGACY模式后,原有的coalesce多格式尝试逻辑即可正常运行,因为每个格式解析失败会返回null,coalesce会依次尝试直到找到匹配的格式。
内容的提问来源于stack exchange,提问作者Vivek Puurkayastha

