【问题标题】:Filtering and reassigning pyspark dataframe in a loop在循环中过滤和重新分配 pyspark 数据帧
【发布时间】:2018-04-30 13:40:44
【问题描述】:

我正在尝试遍历 Pyspark 数据帧的所有列,计算 IQR 以过滤上层异常值,并重新分配相同的数据帧。有 200 多列。这是可行的,但是随着循环的前进,它会变得越来越慢。我怀疑问题可能出在重新分配数据帧(dfBufferOutlier=dfBufferOutlier.filter(col(i)

nameList=dfBuffer.schema.names[1:]    
dfBufferOutlier=dfBuffer
for i in nameList:
    cuenta=dfBufferOutlier.filter(col(i)>0).count()
    if cuenta>0:
        Q1, Q2, Q3 = dfBufferOutlier.filter(col(i)>0).approxQuantile(col=i,probabilities=[0.25,0.5,0.75],relativeError=0.005)
        IQR=(Q3-Q1)
        top_limit=Q3+1.5*IQR
        dfBufferOutlier=dfBufferOutlier.filter(col(i)<top_limit)

肯定有另一种方法可以改善这一点,但不知道如何......请帮忙?

【问题讨论】:

    标签: loops apache-spark dataframe filter pyspark


    【解决方案1】:

    如果你愿意稍微修改一下程序(可以说你还是应该修改它),你可以一次完成所有事情:

    from pyspark.sql.functions import col, count, when
    from functools import reduce
    from operator import and_
    
    def iqr_filter(df, cols, relativeError=0.005):
        # Find quantiles
        quantiles = (df
            # Convert values <= 0 to NULL
            .select([when(col(c) > 0, col(c)).alias(c) for c in cols])
            .approxQuantile(cols, [0.25, 0.5, 0.75], relativeError))
    
        # Compute thresholds
        thresholds = [
            q3 + 1.5 * (q3 - q1) for q1, _, q3 in quantiles
        ]
        # Create SQL expression of form c1 < t1 AND c2 < t2 AND ... AND cn < tn
        expr  = reduce(and_, [col(c) < t for c, t in zip(cols, thresholds)])
        # Filter
        return df.where(expr)
    

    用法:

    iqr_filter(dfBuffer, dfBuffer.schema.names[1:])
    

    相对误差非常低,这可能仍然是一项艰巨的工作,所以我强烈建议在列数增加后稍微放松一下。

    请注意,结果将与您当前方法生成的结果不同:

    • 在您的情况下,每次传递都会计算一组非递增行的分位数,因此它取决于列的顺序。
    • 这会为所有行计算一次分位数,并且不依赖于列的顺序。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-03-16
      • 2021-04-10
      • 2017-10-21
      • 1970-01-01
      • 2022-07-21
      • 2020-08-05
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多