PySpark连接Impala CAST转timestamp报错如何定位异常字段值
问题说明
使用PySpark通过JDBC连接Impala执行查询时触发Timestamp类型转换错误,报错信息如下:
[Cloudera][JDBC](10140) Error converting value to Timestamp.
触发报错的原SQL:
SELECT src_user, CAST(start_time as timestamp) as start_time_ts, start_time, dest_ip, src_ip, count(*) as `count` FROM mytable WHERE start_time like '2022-06%' AND src_ip = '2.3.4.5' AND rule_name like '%XASDF%' group by 1, 2, 3, 4, 5 order by 2 desc)
经测试,移除SQL中的order by子句后查询可正常返回结果,timestamp字段可完成正常转换,初步判断匹配2022-06%前缀的start_time字段中存在无法被正确解析为Timestamp类型的异常值。
当前使用的PySpark读取代码如下:
df = spark.read.jdbc(url= 'jdbc:impala://asdf.com/dasafsda', table = select_sql, properties = properties) df.printSchema() df.show(truncate=False)
需求是配置错误处理逻辑,直接返回触发转换错误的具体字段值,而非仅返回通用错误提示。
解决步骤
- 首先调整JDBC连接参数,关闭驱动侧自动时间类型转换,强制把时间类型字段以字符串形式返回,避免驱动层提前抛出错误导致无法拿到原始值。在原有
properties配置中新增以下参数:# 原有配置项保留,新增以下两项 "mapred.jdbc.timestamp.as.string": "true", "fetch.size": "1000" - 修改查询SQL,移除Impala侧的
CAST(start_time as timestamp)逻辑和order by逻辑,先把原始start_time字符串字段全量拉取到Spark侧再做处理,避免Impala执行阶段提前做类型校验阻断查询:SELECT src_user, start_time, dest_ip, src_ip, count(*) as `count` FROM mytable WHERE start_time like '2022-06%' AND src_ip = '2.3.4.5' AND rule_name like '%XASDF%' group by 1, 2, 3, 4 - 在Spark侧实现时间转换和异常值捕获,直接定位所有异常字段值:
from pyspark.sql import functions as F # 以字符串形式拉取全量数据 df = spark.read.jdbc( url='jdbc:impala://asdf.com/dasafsda', table=select_sql, properties=properties ) # 按实际时间格式做转换,解析失败的记录对应字段会返回null df_with_parse = df.withColumn( "start_time_ts", F.to_timestamp(F.col("start_time"), "yyyy-MM-dd HH:mm:ss") ) # 筛选所有转换失败的异常记录,直接输出异常的start_time值 error_values = df_with_parse.filter( F.col("start_time_ts").isNull() & F.col("start_time").isNotNull() ).select("start_time").distinct() print("触发转换错误的异常字段值:") error_values.show(truncate=False) # 过滤出有效数据后再做排序等后续逻辑 valid_df = df_with_parse.filter(F.col("start_time_ts").isNotNull()).orderBy(F.col("start_time_ts").desc) - 问题原因补充:带
order by就触发报错、移除后暂时正常,是因为Impala执行带排序的查询时,会在结果返回前统一做类型校验和转换,异常值会在拉取数据的第一阶段触发驱动报错;不带排序时驱动是分批拉取数据,若异常值不在首批返回的批次中就不会立刻抛错,但异常值始终存在,上述方案可完整定位所有格式不符合要求的异常值。
内容的提问来源于stack exchange,提问作者mikeb
相关产品推荐
相关产品推荐

