PySpark获取带时区当前时报错:SQLException不支持TIMESTAMP_WITH_TIMEZONE
解决TIMESTAMP_WITH_TIMEZONE不支持的报错问题
核心原因
你的JDBC驱动不支持TIMESTAMP_WITH_TIMEZONE类型,直接查询current_timestamp会返回该类型,导致抛出SQLException。下面是几种可行的解决方式:
方法1:在SQL查询中直接格式化带时区的字符串
用数据库自带的函数把current_timestamp转换成目标格式的字符串,让JDBC返回字符串类型,避开类型不支持的问题。
针对不同数据库的示例:
- MySQL:
select concat(date_format(current_timestamp(), '%Y-%m-%d %H:%i:%s.%f'), ' ', @@session.time_zone) as create_ts, AB.emp_id, AB.emp_name from EMPLOYEE AB - PostgreSQL:
select to_char(current_timestamp, 'YYYY-MM-DD HH24:MI:SS.FF3 TZ') as create_ts, AB.emp_id, AB.emp_name from EMPLOYEE AB - Oracle:
select to_char(current_timestamp, 'YYYY-MM-DD HH24:MI:SS.FF3 TZH:TZM') as create_ts, AB.emp_id, AB.emp_name from EMPLOYEE AB
修改后的Spark代码(以PostgreSQL为例):
val test = source_df.option("query", "select to_char(current_timestamp, 'YYYY-MM-DD HH24:MI:SS.FF3 TZ') as create_ts, AB.emp_id, AB.emp_name from EMPLOYEE AB").load()
方法2:先读取普通TIMESTAMP,再在Spark侧添加时区格式化
第一步:修改SQL返回普通TIMESTAMP
把current_timestamp转换为不带时区的TIMESTAMP类型,确保JDBC能正常读取:
select cast(current_timestamp as timestamp) as create_ts, AB.emp_id, AB.emp_name from EMPLOYEE AB
第二步:Spark中格式化带时区的字符串
导入Spark日期函数,指定目标时区后转换格式:
import org.apache.spark.sql.functions._ // 替换为你需要的时区,比如Asia/Kolkata对应+0530 val target_timezone = "Asia/Kolkata" val formatted_df = test.withColumn( "create_ts", date_format(to_timestamp(col("create_ts")).atTimeZone(target_timezone), "yyyy-MM-dd HH:mm:ss.SSS Z") )
方法3:配置JDBC驱动的类型映射(部分驱动适用)
部分JDBC驱动支持通过URL参数调整类型映射,将TIMESTAMP_WITH_TIMEZONE转为字符串。比如PostgreSQL驱动可以在连接URL中添加:
jdbc:postgresql://host:port/db?stringtype=unspecified
这样驱动会把不支持的类型转为字符串返回,之后你可以在Spark中按需解析格式。
内容的提问来源于stack exchange,提问作者Shini
相关产品推荐
相关产品推荐

