Spark Scala通过JDBC写入PostgreSQL TIME类型字段的解决方案咨询
嘿,这个问题我之前在项目里也碰到过,确实挺头疼的——Spark的类型系统和PostgreSQL的TIME类型衔接起来总有小坑,不过有几个靠谱的办法可以解决,而且不用改数据库表结构,也不会让值变成null:
方案一:直接用Spark的TimeType(推荐,Spark 3.x+适用)
其实Spark是有TimeType这个类型的,你可能没注意到它在org.apache.spark.sql.types包下。我们可以把读取到的hh:mm:ss格式字符串转换成TimeType类型,这样JDBC驱动就能正确映射到PostgreSQL的TIME字段了。
修改你的字段转换逻辑
在原来的类型匹配里加上TimeType的处理:
import org.apache.spark.sql.types.TimeType import java.sql.Time import java.time.LocalTime import java.time.format.DateTimeFormatter val correctValue = field.typ match { case StringType => field.value case IntegerType => val intRegex = """(\d+)""".r field.value match { case intRegex(num) => num.toInt case _ => null } case DoubleType => field.value.toDouble case TimeType => // 把hh:mm:ss字符串转成java.sql.Time val timeFormatter = DateTimeFormatter.ofPattern("HH:mm:ss") val localTime = LocalTime.parse(field.value, timeFormatter) Time.valueOf(localTime) }
更新Schema生成方法
在fileSchemaToStructType里,直接使用TimeType定义对应的StructField:
import org.apache.spark.sql.types.TimeType def fileSchemaToStructType(fileSchema: FileSchema): StructType = { val fields = fileSchema.fields.map { field => StructField( field.name, field.typ, // 这里如果field.typ是TimeType就直接用 nullable = true ) } StructType(fields) }
这样生成的DataFrame里的TimeType列,通过JDBC写入PostgreSQL时,驱动会自动转换成TIME类型,不会报错。
方案二:用StringType配合JDBC自定义类型映射(兼容Spark 2.x)
如果你的Spark版本比较旧(比如2.x),对TimeType支持不够友好,可以用StringType存储hh:mm:ss字符串,然后在写入JDBC时明确指定该列对应PostgreSQL的TIME类型:
写入时指定columnTypes
在调用write.jdbc()的时候,添加columnTypes参数,把对应的字符串列映射成TIME:
val jdbcUrl = "jdbc:postgresql://your-host:5432/your-db" val connectionProps = new java.util.Properties() connectionProps.setProperty("user", "your-user") connectionProps.setProperty("password", "your-pass") df.write .mode(SaveMode.Append) .option("columnTypes", "your_time_column=TIME") // 替换成你的列名 .jdbc(jdbcUrl, "your_table", connectionProps)
PostgreSQL的JDBC驱动会自动把符合hh:mm:ss格式的字符串转换成TIME类型存入数据库。
关于你提到的logical_time_type元数据
那个元数据其实就是Spark的TimeType自带的标记,当你使用TimeType时,Spark会自动为该字段添加logical_time_type: true的元数据,JDBC驱动正是通过这个标记来识别它应该映射到数据库的TIME类型,所以方案一本质上就是利用了这个机制,不用手动去设置元数据。
这两个方案我都试过,都能解决问题,你可以根据自己的Spark版本选择合适的方式~
备注:内容来源于stack exchange,提问作者QueryQuasar

