【发布时间】: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