Glue中Spark SQL带UTC偏移字符串转Timestamp报错及优化咨询
AWS Glue Spark 3.0 时间字符串带时区偏移转换方案
问题场景
在AWS Glue Spark 3.0集群上处理S3 CSV文件,需将格式为20231021134021+0100的时间字符串,转换为叠加时区偏移后的Timestamp(如2023-10-21 14:40:21),最终写入Aurora MySQL。
错误尝试代码
以下代码执行时出现解析错误:
from awsglue.dynamicframe import DynamicFrame # 尝试解析时间并叠加时区偏移 spark_sql = glueContext.sql(""" SELECT DATE_ADD( to_timestamp(SUBSTRING(your_datetime_column, 1, 14), 'yyyyMMddHHmmss'), INTERVAL CAST(SUBSTRING(your_datetime_column, 15, 2) AS INT) HOURS + CAST(SUBSTRING(your_datetime_column, 17, 2) AS INT) MINUTES ) AS converted_time FROM my_data """) result_df = spark.sql(spark_sql) result_dyf = DynamicFrame.fromDF(result_df, glueContext, "result_dyf")
错误原因
Spark SQL的INTERVAL语法不支持直接将多个CAST运算结果相加后作为参数。虽然单独的INTERVAL X HOURS或+ INTERVAL 1 MINUTE是合法语法,但INTERVAL (A+B)的写法不符合规则。
可行解决方案
方案1:拆分INTERVAL分别叠加
将小时和分钟偏移拆分为两个独立的INTERVAL表达式,依次加到基础时间上:
from awsglue.dynamicframe import DynamicFrame spark_sql = glueContext.sql(""" SELECT to_timestamp(SUBSTRING(your_datetime_column, 1, 14), 'yyyyMMddHHmmss') + INTERVAL CAST(SUBSTRING(your_datetime_column, 15, 2) AS INT) HOURS + INTERVAL CAST(SUBSTRING(your_datetime_column, 17, 2) AS INT) MINUTES AS converted_time FROM my_data """) # 注意:glueContext.sql直接返回DataFrame,无需再调用spark.sql result_df = spark_sql result_dyf = DynamicFrame.fromDF(result_df, glueContext, "result_dyf")
方案2:原生解析带时区的时间字符串(最优解)
Spark 3.0支持直接解析带时区偏移的格式,无需手动拆分偏移量。通过正则调整字符串格式后,用to_timestamp自动处理时区:
from awsglue.dynamicframe import DynamicFrame spark_sql = glueContext.sql(""" SELECT to_timestamp( regexp_replace(your_datetime_column, '(\\d{14})([+-]\\d{4})', '$1 $2'), 'yyyyMMddHHmmss Z' ) AS converted_time FROM my_data """) result_df = spark_sql result_dyf = DynamicFrame.fromDF(result_df, glueContext, "result_dyf")
说明:用
regexp_replace把20231021134021+0100转为20231021134021 +0100,匹配Spark支持的Z时区格式,to_timestamp会自动叠加偏移量转换为Session默认时区的Timestamp。
方案3:计算总偏移分钟数叠加
将小时和分钟转换为总分钟数,通过date_add或直接数值相加实现:
from awsglue.dynamicframe import DynamicFrame spark_sql = glueContext.sql(""" SELECT to_timestamp(SUBSTRING(your_datetime_column, 1, 14), 'yyyyMMddHHmmss') + INTERVAL (CAST(SUBSTRING(your_datetime_column, 15, 2) AS INT)*60 + CAST(SUBSTRING(your_datetime_column, 17, 2) AS INT)) MINUTES AS converted_time FROM my_data """) result_df = spark_sql result_dyf = DynamicFrame.fromDF(result_df, glueContext, "result_dyf")
内容的提问来源于stack exchange,提问作者Santhosh
相关产品推荐
相关产品推荐

