Spark Scala中如何过滤含ERROR的Struct类型列?
解决Spark Scala中Struct类型列的ERROR过滤问题
因为rlike仅支持字符串类型,无法直接作用于Struct列,你可以通过以下两种可行方式实现需求:
方法一:直接遍历Struct的子字段匹配
明确指定Struct中的两个字符串字段,分别判断是否包含ERROR,用逻辑或连接条件,这种方式性能最优且精准:
import org.apache.spark.sql.functions.col val errorJobRunsDF = df.filter( col("state.life_cycle_state").rlike("ERROR") || col("state.state_message").rlike("ERROR") ).select("JobID")
方法二:将Struct转为字符串后统一匹配
如果后续Struct可能新增字段,不想每次修改过滤条件,可以把Struct转为JSON字符串,再用rlike匹配ERROR,这种方式更灵活:
import org.apache.spark.sql.functions.{col, to_json} val errorJobRunsDF = df.filter( to_json(col("state")).rlike("ERROR") ).select("JobID")
注意:这种方式会引入JSON转换的轻微性能开销,且需确保Struct转为JSON后,ERROR不会被JSON的格式符号干扰(比如字段名中不会出现ERROR)。
内容的提问来源于stack exchange,提问作者Santosh Reddy Kommidi
相关产品推荐
相关产品推荐

