Spark Streaming解析Event Hub消息JSON拆分多列写入Delta表报错
问题解决方案
错误原因
你遇到的cannot resolve 'body.*' given input columns 'body'报错核心原因是:
调用from_json生成的Struct类型列没有命名为body,df4的结构仅包含一个默认名称为from_json(body, <schema>)的Struct列,不存在名为body的Struct字段,因此无法执行body.*展开操作。
修复步骤
- 调整列别名设置,保证
from_json返回的Struct列正确命名为body - 优化中间步骤的冗余代码,避免无效转换
- 提前确认
jsonSchema与Event Hub传输的JSON结构完全匹配,避免展开后出现全空列
修复后完整代码
val incomingStream = spark.readStream.format("eventhubs").options(customEventhubParameters.toMap).load() // 直接转换body为string并保留列名 val df = incomingStream.select($"body".cast(StringType).alias("body")) // from_json处理后给返回的Struct列显式起别名body val df4 = df.select(from_json(col("body"), jsonSchema).alias("body")) // 此时df4存在名为body的Struct列,可以正常展开 val df5 = df4.select("body.*") df5.writeStream .format("delta") .outputMode("append") .option("ignoreChanges", "true") .option("checkpointLocation", "/mnt/abc/checkpoints/samplestream") .start("/mnt/abc/samplestream")
校验方法
生成df4后可调用df4.printSchema()确认列结构,正常输出会包含名为body的Struct类型列,其下嵌套的字段就是你定义的jsonSchema中的字段,此时再执行展开操作就不会报错。
内容的提问来源于stack exchange,提问作者NgD
相关产品推荐
相关产品推荐

