PySpark:基于Status列fail条件统计数组中指定元素数量
按对应条件统计数组元素数量
问题场景
现有Spark DataFrame包含两个数组列Status和Description,两个数组的元素位置一一对应。需求是:仅统计**Status列对应位置值为fail**时,Description列对应位置是"Clear BMC SEL"或"Check SEL"的元素总数,预期结果为2。
示例数据构造代码:
df = spark.createDataFrame(sc.parallelize([ [ ['pass', 'pass', 'fail', 'fail', 'fail'], ['Clear BMC SEL', 'Check SEL', 'Clear BMC SEL', 'CPU', 'Check SEL'] ] ]), ['Status', 'Description'])
原代码仅过滤了Description列的目标值,未关联Status列的对应位置,导致统计结果为4,不符合需求:
df = df.selectExpr('*', 'filter(Description, x -> x = "Clear BMC SEL" or x = "Check SEL") as pass_array') df = df.selectExpr('*', 'size(pass_array) as testFailCount').drop('pass_array')
解决方案
使用arrays_zip将Status和Description两个数组按位置打包成结构体数组,再过滤同时满足以下两个条件的结构体:
Status字段值为failDescription字段值是"Clear BMC SEL"或"Check SEL"
最后统计过滤后的数组长度即可得到目标结果,提供两种实现方式:
方式1:Spark SQL表达式实现
df = df.selectExpr( '*', 'size(filter(arrays_zip(Status, Description), item -> item.Status = "fail" and (item.Description = "Clear BMC SEL" or item.Description = "Check SEL"))) as testFailCount' )
方式2:PySpark API实现(更直观)
from pyspark.sql.functions import arrays_zip, col, filter, size df = df.withColumn( 'zipped_arr', arrays_zip(col('Status'), col('Description')) ).withColumn( 'testFailCount', size( filter( col('zipped_arr'), lambda x: (x.Status == 'fail') & ((x.Description == 'Clear BMC SEL') | (x.Description == 'Check SEL')) ) ) ).drop('zipped_arr')
执行后结果符合预期:
+--------------------+--------------------+-------------+ | Status| Description|testFailCount| +--------------------+--------------------+-------------+ |[pass, pass, fail...|[Clear BMC SEL, C...| 2| +--------------------+--------------------+-------------+
内容的提问来源于stack exchange,提问作者John Stud
相关产品推荐
相关产品推荐

