【问题标题】:PySpark - an efficient way to find DataFrame columns with more than 1 distinct valuePySpark - 一种查找具有超过 1 个不同值的 DataFrame 列的有效方法
【发布时间】:2019-09-03 02:52:12
【问题描述】:

我需要一种有效的方法来列出和删除 Spark DataFrame 中的一元列(我使用 PySpark API)。我将一元列定义为最多具有一个不同值的列,出于定义的目的,我也将null 计为一个值。这意味着在某些行中具有不同 non-null 值而在其他行中具有不同 null 值的列不是一元列。

根据this question 的答案,我设法编写了一种有效的方法来获取空列列表(它们是我的一元列的子集)并将它们删除,如下所示:

counts = df.summary("count").collect()[0].asDict()
null_cols = [c for c in counts.keys() if counts[c] == '0']
df2 = df.drop(*null_cols)

基于我对 Spark 内部工作原理的非常有限的理解,这很快,因为方法摘要同时操作整个数据帧(我的初始数据帧中大约有 300 列)。不幸的是,我找不到类似的方法来处理第二种类型的一元列——那些没有null 值但为lit(something) 的列。

我目前拥有的是这个(使用我从上面的代码 sn-p 获得的df2):

prox_counts = (df2.agg(*(F.approx_count_distinct(F.col(c)).alias(c)
                         for c in df2.columns
                         )
                       )
                  .collect()[0].asDict()
               )
poss_unarcols = [k for k in prox_counts.keys() if prox_counts[k] < 3]
unar_cols = [c for c in poss_unarcols if df2.select(c).distinct().count() < 2]

基本上,我首先以快速但近似的方式找到可能是一元的列,然后更详细、更缓慢地查看“候选”。

我不喜欢它的是 a) 即使进行了近似的预选,它仍然相当慢,即使此时我只有大约 70 列(大约 600 万列)也需要一分钟多的时间才能运行行)和b)我使用approx_count_distinct和神奇的常数3approx_count_distinct不算null,因此3而不是2)。由于我不确定approx_count_distinct 在内部是如何工作的,我有点担心3 不是一个特别好的常数,因为该函数可能会估计不同(non-null)值的数量,比如5是 1,因此可能需要更高的常数来保证候选列表 poss_unarcols 中没有任何遗漏。

有没有更聪明的方法来做到这一点,理想情况下,我什至不必单独删除空列并一口气完成所有操作(尽管这实际上非常快,所以这是一个大问题) ?

【问题讨论】:

    标签: python apache-spark dataframe pyspark


    【解决方案1】:

    我建议你看看下面的函数

    pyspark.sql.functions.collect_set(col)
    

    https://spark.apache.org/docs/latest/api/python/pyspark.sql.html?highlight=dataframe

    它将返回 col 中的所有值,并消除相乘的元素。然后您可以检查结果的长度(是否等于 1)。我想知道性能,但我认为它肯定会击败 distinct().count()。让我们星期一看看:)

    【讨论】:

      【解决方案2】:

      您可以 df.na.fill("some non exisitng value").summary() 然后从原始数据框中删除相关列

      【讨论】:

      • 这个问题是列有不同的数据类型,所以我不得不多次调用它。另外,我意识到我对一元的定义意味着一旦我删除了完全为空的列,我只需要检查其中没有空值的列 - 所有其他列至少有 1 行等于 null 和 1 行等于某事else 所以它们不是一元的。
      【解决方案3】:

      到目前为止,我找到的最佳解决方案是这个(它比其他建议的答案更快,虽然不理想,见下文):

      rows = df.count()
      nullcounts = df.summary("count").collect()[0].asDict()
      del nullcounts['summary']
      nullcounts = {key: (rows-int(value)) for (key, value) in nullcounts.items()}
      
      # a list for columns with just null values
      null_cols = []
      # a list for columns with no null values
      full_cols = []
      
      for key, value in nullcounts.items():
          if value == rows:
              null_cols.append(key)
          elif value == 0:
              full_cols.append(key)
      
      df = df.drop(*null_cols)
      
      # only columns in full_cols can be unary
      # all other remaining columns have at least 1 null and 1 non-null value
      try:
          unarcounts = (df.agg(*(F.countDistinct(F.col(c)).alias(c) for c in full_cols))
                          .collect()[0]
                          .asDict()
                        )
          unar_cols = [key for key in unarcounts.keys() if unarcounts[key] == 1]
      except AssertionError:
          unar_cols = []
      
      df = df.drop(*unar_cols)
      

      这工作相当快,主要是因为我没有太多的“完整列”,即不包含 null 行的列,我只遍历所有这些行,使用快速 summary("count") 方法进行分类尽可能多的列。

      遍历一列的所有行对我来说似乎非常浪费,因为一旦找到两个不同的值,我就真的不在乎该列的其余部分是什么。我不认为这可以在 pySpark 中解决(但我是初学者),这似乎需要一个 UDF,而 pySpark UDF 太慢了,它不可能比使用countDistinct() 更快。尽管如此,只要数据框中有很多列没有null 行,这种方法就会很慢(而且我不确定有多少人可以信任approx_count_distinct() 来区分一列中的一个或两个不同的值)

      据我所知,它胜过collect_set() 方法并且填充null 值实际上是没有必要的,因为我意识到(请参阅代码中的 cmets)。

      【讨论】:

        猜你喜欢
        • 2022-08-19
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-12-25
        • 2020-05-04
        相关资源
        最近更新 更多