You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 17:36:00