将Databricks流水线事件日志查询SQL转换为PySpark代码求助
Databricks事件日志查询PySpark改写方案
你的写法失效核心原因是**details列是字符串类型**,不是结构体/嵌套类型,所以没法直接用SQL的冒号语法或DataFrame的点语法访问嵌套字段,必须先把字符串解析为JSON格式再提取目标字段。以下两种方案都可以解决问题:
方案1:用get_json_object快速提取字段
适合只需要提取少量指定字段的场景,无需定义完整JSON schema:
from pyspark.sql.functions import get_json_object # 过滤后提取指定JSON路径的字段,用alias重命名匹配SQL结果 result_df = df.filter(df.event_type == 'user_action') \ .select( "timestamp", get_json_object(df.details, "$.user_action.action").alias("action"), get_json_object(df.details, "$.user_action.user_name").alias("user_name") )
- 说明:
get_json_object的第二个参数是JSONPath表达式,$代表JSON根节点,后续按层级写路径即可定位到目标字段。
方案2:用from_json解析为结构体(推荐复用场景)
如果后续需要多次访问details里的嵌套字段,建议先把字符串解析为结构体类型,之后就能像操作普通嵌套列一样用点语法访问:
from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StructField, StringType # 定义与details字符串对应的JSON结构schema details_schema = StructType([ StructField("user_action", StructType([ StructField("action", StringType()), StructField("user_name", StringType()) ])) ]) # 先解析字符串为结构体,再提取字段 result_df = df.filter(df.event_type == 'user_action') \ .withColumn("details_parsed", from_json(df.details, details_schema)) \ .select( "timestamp", "details_parsed.user_action.action", "details_parsed.user_action.user_name" )
内容的提问来源于stack exchange,提问作者Lilybone
相关产品推荐
相关产品推荐

