【发布时间】:2020-06-03 18:14:07
【问题描述】:
我现在有一个大型 python 项目,其中驱动程序有一个函数,该函数使用 for 循环遍历我的 GCP(谷歌云平台)存储桶上的每个文件。我正在使用 CLI 将作业提交到 GCP 并让作业在 GCP 上运行。
对于在此 for 循环中遍历的每个文件,我正在调用一个函数 parse_file(...) 来解析文件并调用处理此文件的其他函数的序列。
整个项目运行需要几分钟,速度很慢,而且驱动程序还没有使用太多的 PySpark。问题是该文件级 for 循环中的每个 parse_file(...) 都按顺序执行。是否可以使用 PySpark 并行化该文件级 for 循环以对所有这些文件并行运行 parse_file(...) 函数以减少程序执行时间并提高效率?如果是这样,由于程序没有使用 PySpark,是否需要进行大量代码修改才能使其并行化?
所以程序的功能是这样的
# ... some other codes
attributes_table = ....
for obj in gcp_bucket.objects(path):
if obj.key.endswith('sys_data.txt'):
#....some other codes
file_data = (d for d in obj.download().decode('utf-8').split('\n'))
parse_file(file_data, attributes_table)
#....some other codes ....
如何使用 PySpark 并行化这部分,而不是一次使用 for 循环遍历文件?
【问题讨论】:
标签: apache-spark for-loop pyspark parallel-processing