PySpark提取MongoDB生日字段数值问题求助
问题分析与解决方案
核心问题原因
你的代码中CASE WHEN的判断逻辑存在缺陷:Spark的get_json_object当指定路径不存在时,返回的是空字符串""而非NULL,因此原判断条件IS NOT NULL无法触发后续分支。此外,部分分支未显式转换数值类型,也可能导致结果不符合预期。
修正后的代码
from pyspark.sql.functions import expr df_extracted = df.withColumn( 'birthday_extracted', expr(""" CASE -- 处理直接存储sec的结构,先判断非空再转成BIGINT WHEN trim(get_json_object(birthday, '$.sec')) != '' THEN TRY_CAST(get_json_object(birthday, '$.sec') AS BIGINT) -- 处理sec嵌套$numberLong的结构 WHEN trim(get_json_object(birthday, '$.sec.$numberLong')) != '' THEN TRY_CAST(get_json_object(birthday, '$.sec.$numberLong') AS BIGINT) -- 处理date字段(毫秒级转秒级) WHEN trim(get_json_object(birthday, '$.date')) != '' THEN TRY_CAST(TRY_CAST(get_json_object(birthday, '$.date') AS DOUBLE) / 1000 AS BIGINT) -- 所有未匹配的情况返回NULL ELSE NULL END """) )
关键修改点说明
- 判断条件调整:用
trim(...) != ''替代IS NOT NULL,适配get_json_object返回空字符串的特性,避免遗漏有效分支。 - 类型安全转换:使用
TRY_CAST替代普通CAST,如果提取的字符串无法转换为数值时,会返回NULL而非触发报错,保证任务稳定性。 - 统一输出类型:所有分支最终都转换为
BIGINT,确保birthday_extracted字段类型一致。
调试建议
如果仍有分支不生效,先运行以下代码查看各路径的提取结果,确认JSON结构与路径是否匹配:
from pyspark.sql.functions import get_json_object df.select( "birthday", get_json_object("birthday", "$.sec").alias("sec_value"), get_json_object("birthday", "$.sec.$numberLong").alias("sec_numberLong_value"), get_json_object("birthday", "$.date").alias("date_value") ).show(truncate=False)
内容的提问来源于stack exchange,提问作者Padam Rai
相关产品推荐
相关产品推荐

