【问题标题】:Spark Partitions: Loading a file from the local file system on a Single Node ClusterSpark Partitions:从单节点集群上的本地文件系统加载文件
【发布时间】:2018-07-29 00:02:33
【问题描述】:

我有兴趣了解 Spark 在从本地文件系统加载文件时如何创建分区。

我正在使用 Databricks 社区版来学习 Spark。当我使用 sc.textfile 命令加载一个只有几千字节(大约 300 kb)的文件时,spark 默认会创建 2 个分区(由 partitions.length 给出)。当我加载大约 500 MB 的文件时,它会创建 8 个分区(等于机器中的内核数)。

enter image description here

这里的逻辑是什么?

另外,我从文档中了解到,如果我们从本地文件系统加载并使用集群,则该文件必须位于属于该集群的所有机器上的相同位置。这不会创建重复项吗? Spark 如何处理这种情况?如果您能指出阐明这一点的文章,那将有很大帮助。

谢谢!

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    当 Spark 从本地文件系统读取时,默认的分区数(由 defaultParallelism 标识)是所有可用内核数

    sc.textFile 将分区数计算为 defaultParallelism(本地 FS 的可用内核)和 2 之间的最小值。

    def defaultMinPartitions: Int = math.min(defaultParallelism, 2)
    

    引用自:spark code

    第一种情况:文件大小 - 300KB

    分区数计算为 2,因为文件大小非常小。

    在第二种情况下:文件大小 - 500MB

    分区数等于默认并行度。在你的情况下,它是 8。

    从 HDFS 读取时,sc.textFile 将取 minPartitions 和根据 hadoop 输入拆分大小除以块大小计算的拆分数量之间的最大值。

    但是,当使用带有压缩文件(file.txt.gz 不是 file.txt 或类似文件)的 textFile 时,Spark 会禁用拆分,导致 RDD 只有 1 个分区(因为对 gzip 文件的读取无法并行化)。

    关于从集群中的本地路径读取数据的第二个查询:

    文件需要在集群中的所有机器上都可用,因为 Spark 可能会在集群中的机器上启动执行器,执行器将使用 (file://) 读取文件。

    为了避免将文件复制到所有机器上,如果您的数据已经在 NFS、AFS 和 MapR 的 NFS 层等网络文件系统之一中,那么您只需指定一个文件即可将其用作输入:/ / 小路;只要文件系统安装在每个节点上的相同路径,Spark 就会处理它。每个节点都需要有相同的路径。 请参考:https://community.hortonworks.com/questions/38482/loading-local-file-to-apache-spark.html

    【讨论】:

    • 很棒的解释§
    • 感谢@Lakshman-Battini 的解释。这是我的疑问。基于 Math.min(defaultparallelism,2),在这两种情况下(文件大小 300kb 和 500MB),分区数不应该是 2 吗?因为那将是最小值? 500 MB 文件的 8 位如何?至于第二个问题,我知道文件应该位于同一路径中属于集群的所有机器上。这是否意味着每个执行程序都会加载整个文件?如果他们每个人都创建分区,Spark 如何知道避免重复?
    • 这是文件的可拆分性。在 Hadoop 中,如果文件是可拆分的,spark 将启动并行任务来读取每个块(块大小为 64 / 129 MB)。即使在本地 FS 的情况下,Spark 也足够智能,如果文件是可拆分的,Spark 会启动并行读取文件的任务。
    • 对于您的第二次查询,每个执行程序不会加载整个文件,它只会从文件块的开头和结尾读取。
    猜你喜欢
    • 1970-01-01
    • 2018-01-25
    • 2018-11-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多