PySpark使用from_json转换Kafka JSON至StructType时遇属性错误
问题解决:Spark from_json 报错 'StructField' object has no attribute '_get_object_id'
错误原因
你在from_json的第二个参数位置直接用when返回schema_session_start(StructType对象),但when要求返回的是Column类型,而StructType不属于Column,因此触发了属性错误。另外你的数据示例存在JSON格式错误,Schema定义也有一处不匹配。
具体修正步骤
1. 修正Schema定义
数据中FN字段是布尔值true,但你定义为StringType,需要改成BooleanType:
schema_session_start = StructType([ StructField("ID", StringType()), StructField("SID", StringType()), StructField("EP", LongType()), StructField("IP", StringType()), StructField("LN", StringType()), StructField("VN", StringType()), StructField("DV", StructType([ StructField("MK", StringType()), StructField("MDL", StringType()), StructField("OS", StringType()), StructField("OSVN", StringType()), StructField("AR", StringType()) ])), StructField("MC", StringType()), StructField("FN", BooleanType()), # 修正为BooleanType StructField("NW", StructType([ StructField("TP", StringType()) ])), StructField("AL", StringType()), StructField("EN", StringType()) ])
2. 修正JSON数据格式
数据示例中NW字段语法错误,原写法"NW":"TP":"wifi_5"}不符合JSON规范,应改为"NW":{"TP":"wifi_5"},修正后的完整JSON:
{"ID":"651551912131b2.07017577","SID":"169156360280217644","EP":1695895952305,"IP":"10.10.10.10","LN":"","VN":"2.4.0.0","DV":{"MK":"Jio","MDL":"JHSD200","OS":"JioOS 2","OSVN":"9","AR":"armeabi-v7a"},"MC":"02:00:00:00:00:00","FN":true,"NW":{"TP":"wifi_5"},"AL":"GRIPdemo","EN":"Session_Start"}
3. 修正转换代码
from_json的第二个参数必须是固定Schema,条件判断应该包裹整个from_json调用,或者先过滤数据:
方式一:只处理符合条件的行(过滤后转换)
filtered_df = condition_df.filter(col("EN") == "Session_Start") array_df = filtered_df.withColumn("value_json", from_json(col("value"), schema_session_start))
方式二:保留所有行,不符合条件的返回null
array_df = condition_df.withColumn( "value_json", when( col("EN") == "Session_Start", from_json(col("value"), schema_session_start) ).otherwise(lit(None)) )
内容的提问来源于stack exchange,提问作者Virendar Kumar
相关产品推荐
相关产品推荐

