【问题标题】:Processing a large file with pyspark locally在本地使用 pyspark 处理大文件
【发布时间】:2019-08-07 09:40:11
【问题描述】:

我是 PySpark 的新手,只是用它来处理数据。

我有一个 120GB 的文件,其中包含超过 10.5 亿行。我能够对文件进行聚合和过滤,并使用 coalesce() 函数将结果输出到 CSV 文件,没有任何问题。

我的挑战是,当我尝试通读文件中的每一行以执行一些计算时,我的 spark 作业使用 .collect() 或 .toLocalIterator() 函数失败。当我限制读取的行数时,它工作正常。

请问,我该如何解决这个挑战?是否可以按位读取行,例如一次一行还是一次一大块?

我在 64GB RAM 的计算机上本地运行 Spark。

下面是我的 Python 代码示例:

sql = "select * from table limit 1000"
details = sparkSession.sql(sql).collect()
for detail in details:
    #do some computation

下面是我失败的python代码示例:

sql = "select * from table"
details = sparkSession.sql(sql).collect()
for detail in details:
    #do some computation

这是我提交 Spark 作业的方式

spark-submit --driver-memory 16G --executor-memory 16G python_file.py

非常感谢。

【问题讨论】:

  • 收集数据就像首先将整个数据集加载到内存中一样好。如果您考虑到所有重复,甚至更糟。
  • 我不明白为什么您必须使用collect() 进行处理。您可以对每一行应用转换。使用collect(),您将所有内容都放入驱动程序,即将所有内容加载到内存中。这显然是错误的,在这种情况下,你完全可以不使用 Spark。通过 Spark 正确执行或不使用 Spark。但不像您在上面所做的那样......任何map() 进程显然是要走的路。您应该再次阅读有关 Spark 以及文档的信息

标签: python apache-spark pyspark


【解决方案1】:

您解决问题的方法是错误的。 collect 方法将完整文件(由于反序列化实际上可能需要超过 120GB)加载到驱动程序内存(单个 pyspark 进程)中,导致内存不足。
根据经验,如果您使用 collect() 方法在 spark 代码中它不好,应该改变。

如果使用得当,spark 将一次只读取部分输入数据(输入拆分)来处理并产生(小得多)存储在执行程序内存中的中间结果。因此它(取决于处理类型)可以处理 16GB 内存的 120GB 文件。

【讨论】:

    猜你喜欢
    • 2015-05-01
    • 2016-09-02
    • 1970-01-01
    • 2011-05-16
    • 1970-01-01
    • 2012-12-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多