【问题标题】:Parallelization in counting Spark dataframe groups in pyspark在 pyspark 中计算 Spark 数据帧组的并行化
【发布时间】:2017-10-12 07:19:14
【问题描述】:

我在名为part-0001part-0002 等的 Linux 机器上的单个目录中有大约 200 个文件。每个都有大约一百万行具有相同的列(称它们为“a”、“b”等)。让 'a','b' 对成为每一行的键(有许多重复项)。

同时,我建立了一个 Spark 2.2.0 集群,有一个 master 和两个 slave,总共有 42 个可用内核。地址是spark://XXX.YYY.com:7077

然后我使用 PySpark 连接到集群并计算每个唯一对的 200 个文件的计数,如下所示。

from pyspark import SparkContext
from pyspark.sql import SQLContext
import pandas as pd

sc = SparkContext("spark://XXX.YYY.com:7077")
sqlContext = SQLContext(sc)

data_path = "/location/to/my/data/part-*"
sparkdf = sqlContext.read.csv(path=data_path, header=True)
dfgrouped = sparkdf.groupBy(['a','b'])
counts_by_group = dfgrouped.count()

我可以看到 Spark 正在处理一系列消息,并且它确实返回了看起来合理的结果。

问题:在执行此计算时,top 没有显示任何证据表明从属内核正在执行任何操作。似乎没有任何并行化。每个从属设备在作业之前都有一个相关的 Java 进程(加上来自其他用户的进程和后台系统进程)。所以看起来主人正在做所有的工作。鉴于有 200 个奇怪的文件,我曾预计在每台从属计算机上运行 21 个进程,直到事情结束(这 当我显式调用 parallelize 时看到的,如下 count = sc.parallelize(c=range(1, niters + 1), numSlices=ncores).map(f).reduce(add)一个单独的实现)。

问题:如何确保 Spark 实际并行化计数?我希望每个核心抓取一个或多个文件,对它在文件中看到的对执行计数,然后将单个结果缩减为单个 DataFrame。我不应该在顶部看到这个吗?我需要显式调用并行化吗?

(FWIW,我见过使用分区的示例,但我的理解是,这是用于在 single 文件的块上分配处理。我的情况是我有很多文件。)

提前致谢。

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    TL;DR可能您的部署没有任何问题。

    我原本预计会看到 21 个进程在运行

    除非您专门将 Spark 配置为每个执行程序 JVM 使用一个内核,否则没有理由发生这种情况。与您在问题中提到的RDD 示例不同,DataFrame API 根本不使用 Python 工作者,Python UserDefinedFunctions 除外。

    同时,JVM 执行器使用线程而不是成熟的系统进程(PySpark 使用后者来避免GIL)。此外,独立模式下的默认spark.executor.cores 等于the available cores on the worker 的数量。因此,无需额外配置,您应该会看到两个执行程序 JVM,每个执行程序使用 21 个数据处理线程。

    总的来说,你应该检查 Spark UI,如果你看到分配给执行者的任务,一切都应该没问题。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-09-18
      • 2019-05-29
      • 1970-01-01
      • 2016-02-27
      • 2023-04-11
      • 2020-02-14
      相关资源
      最近更新 更多