Spark的current_timestamp值写入PostgreSQL带时区timestamp列JDBC方案咨询
问题原因
- 你使用的普通三引号字符串没有添加Scala字符串插值的
s前缀,SQL语句中的$start_time等变量没有被替换为实际值,直接把$start_time字符串传入了PostgreSQL,导致数据库无法识别为合法的带时区时间戳格式。 - 直接拼接SQL的写法存在SQL注入风险,且时间类型转字符串拼接容易出现时区、格式不兼容问题,不推荐使用。
解决方案
使用JDBC的PreparedStatement预编译语句实现参数绑定,既能自动处理时间类型的格式转换,也能避免SQL注入问题,同时调整字符串插值的错误写法、优化连接创建逻辑(不要在循环内重复创建数据库连接),修改后代码如下:
import java.sql.{Connection, DriverManager, PreparedStatement, Timestamp} import org.apache.spark.sql.functions.current_timestamp import spark.implicits._ // 数据处理逻辑保持不变 val df1 = Seq(("Process Name", "Process Description")) .toDF("process_nm", "process_desc") val df2 = df1.withColumn("start_time", current_timestamp) df2.show(false) df2.printSchema // 数据库连接放到循环外创建,避免重复创建连接消耗性能 val db_conn_string = s"jdbc:$db_type://$db_host:$db_port/$db_database" var direct_conn: Connection = null var pstmt: PreparedStatement = null try { direct_conn = DriverManager.getConnection(db_conn_string, db_user, db_pass) // 用?作为参数占位符定义预编译SQL val insertSql = "insert into test_log(process_nm, start_time, process_desc) values (?, ?, ?)" pstmt = direct_conn.prepareStatement(insertSql) df2.collect().foreach(row => { val process_name: String = row.getString(0) val process_description: String = row.getString(1) val start_time: Timestamp = row.getTimestamp(2) println("Before calling-" + process_name + " " + process_description + " " + start_time) // 按位置绑定参数,JDBC驱动会自动适配PostgreSQL的timestamp with timezone类型 pstmt.setString(1, process_name) pstmt.setTimestamp(2, start_time) pstmt.setString(3, process_description) val result = pstmt.executeUpdate() println("inserted-" + result) }) } finally { // 资源释放,避免连接泄漏 if (pstmt != null) pstmt.close() if (direct_conn != null) direct_conn.close() }
补充说明
如果需要明确指定时区适配timestamp with timezone类型,可以使用setTimestamp的重载方法传入时区对应的Calendar实例:
import java.util.Calendar import java.util.TimeZone // 替换上面的setTimestamp行,指定你需要的时区,比如UTC pstmt.setTimestamp(2, start_time, Calendar.getInstance(TimeZone.getTimeZone("UTC")))
内容的提问来源于stack exchange,提问作者saikiran
相关产品推荐
相关产品推荐

