【问题标题】:Checking DataFrame has records in PySpark检查 DataFrame 在 PySpark 中有记录
【发布时间】:2018-07-02 07:11:05
【问题描述】:

这是我第一次使用 Python 或 Spark,我是一名 Java 开发人员。所以我不知道这里最好的解决方法是什么。

我正在与:

  • 为 Hadoop 2.7.3 构建的 Spark 2.2.0
  • Python 2.7.12

我有一个 PySpark 脚本,该脚本执行不同的查询并创建临时视图,直到它最终使用/加入不同的临时视图执行最终查询。它将使用最终执行查询的结果写入文件。

脚本运行良好,但我们发现,当没有数据时,它仍然会创建 200 个文件(全部为空)。我们想在调用 write 方法甚至创建临时视图之前验证它是否确实有数据,所以我们尝试使用if df.count() == 0:,如果是会引发错误,否则,请继续。

我刚刚在执行临时视图之前将验证添加到最后两个数据帧,因此它会尽快中断该过程,并且在执行下一个查询之前。

然后我们在某个地方读到,count 是一种非常昂贵的方法来验证是否有数据,因为它经过所有的刽子手,所以在尝试之前,我们在几个地方更改了推荐的方法:使用df.take(1)df.head(1)df.first(1)。我们最终选择了head(1)

但是,这将执行经过的时间从 30 分钟更改为实际超过 1 小时 40 分钟。

我想知道在不增加太多计算时间的情况下,我可以通过哪种其他方式避免 spark 写入空文件。

由于我是新手,所以我愿意接受建议。

编辑

我已经阅读了这个帖子:How to check if spark dataframe is empty。从这个线程中,我认为我应该使用len(df.head(1)) == 0,这将计算时间从 30 分钟增加到 1h 40m+。

【问题讨论】:

  • 嗯,这正是我上周阅读的主题之一,从那个确切的主题来看,我选择了df.head(1)它减慢了进程的速度远不止一个小时.
  • 我们需要一个MVCE 来复制它,否则它是不可挽救的,我投票结束这个问题。 count 比 head 更昂贵,因为它会扫描所有数据,但在需要实际计算之前你可能正在做其他事情,这可能不是你的瓶颈......
  • 您的某一列的计算成本一定很高。是否有非计算列可以选择:len(df.select('non-computed column').head(1)) == 0
  • 除非他从 jdbc 或类似@Jaco 的东西中提取数据,这就是我们需要 MVCE 的原因...

标签: apache-spark pyspark apache-spark-sql spark-dataframe pyspark-sql


【解决方案1】:

只需获取数据框的 rdd 并检查它是否为空:

df.rdd.isEmpty()

spark 中有两种类型的操作:动作和转换。 Spark 中的所有转换都是惰性的,它们不会立即计算结果。仅在执行操作时才计算转换。动作代价高昂,因为 spark 需要运行所有转换到该点才能运行动作。

【讨论】:

    【解决方案2】:

    @Jaco 我终于做了类似if df.select('my_no_computed_column').head() is None: 的事情,因为显然,没有参数的 head() 将假定为 1 并且根据 Spark 的代码:

        @ignore_unicode_prefix
        @since(1.3)
        def head(self, n=None):
            """Returns the first ``n`` rows.
    
            .. note:: This method should only be used if the resulting array is expected
                to be small, as all the data is loaded into the driver's memory.
    
            :param n: int, default 1. Number of rows to return.
            :return: If n is greater than 1, return a list of :class:`Row`.
                If n is 1, return a single Row.
    
            >>> df.head()
            Row(age=2, name=u'Alice')
            >>> df.head(1)
            [Row(age=2, name=u'Alice')]
            """
            if n is None:
                rs = self.head(1)
                return rs[0] if rs else None
            return self.take(n)
    

    如果没有行,它将返回 None (虽然我可能读错了,我已经用 Java 编程超过 10 年了,Python 和 Spark 对我来说太新了,Python对我来说太奇怪了)。

    它确实大大减少了运行时间。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-11-04
      • 1970-01-01
      • 1970-01-01
      • 2021-05-14
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多