【问题标题】:All the executors are not being used when reading JSON(zipped .gz) in GCP from google dataproc spark cluster using spark-submit使用 spark-submit 从 google dataproc spark 集群中读取 GCP 中的 JSON(zipped .gz) 时,所有执行程序都没有被使用
【发布时间】:2019-10-27 12:35:57
【问题描述】:

我刚刚使用 GCP(dataproc) 和 pyspark 了解了这个美妙的大数据和云技术世界。我有 ~5 GB 大小的 JSON 文件(压缩,gz 文件),其中包含 ~500 万 条记录,我只需要读取每一行并处理满足特定条件的那些行。我有我的工作代码,我发出了一个带有 --num-partitions=5 的 spark-submit,但仍然只有一名工人用于执行该操作。

这是我正在使用的 spark-submit 命令:

spark-submit --num-executors 5 --py-files /home/user/code/dist/package-0.1-py3.6.egg job.py

job.py:

path = "gs://dataproc-bucket/json-files/data_5M.json.gz"
mi = spark.read.json(path)
inf_rel = mi.select(mi.client_id,
                    mi.user_id,
                    mi.first_date,
                    F.hour(mi.first_date).alias('hour'),
                    mi.notes).rdd.map(foo).filter(lambda x: x)
inf_relevance = inf_rel.map(lambda l: Row(**dict(l))).toDF()
save_path = "gs://dataproc-bucket/json-files/output_5M.json"
inf_relevance.write.mode('append').json(save_path)
print("END!!")

Dataproc 配置: (我现在使用免费帐户,一旦我得到工作解决方案将添加更多核心和执行器)

(Debian 9、Hadoop 2.9、Spark 2.4) 主节点:2 个 vCPU,7.50 GB 内存 主磁盘大小:32 GB 5 个工作节点:1 个 vCPU,3.75 GB 内存 主磁盘类型:32 GB

在 spark-submit 之后,我可以在 web UI 中看到添加了 5 个执行程序,但只有 1 个执行程序保持活动状态并执行所有任务,其余 4 个被释放。

我进行了研究,大部分问题都是关于通过 JDBC 访问数据的。

请提出我在这里缺少的内容。

附:最终我会读取 64 个 5 GB 的 json 文件,因此可能会使用 8 个核心 * 100 个工作人员。

【问题讨论】:

  • 你的数据框有多少个分区?检查df.rdd.getNumPartitions() ...我没有使用这个google dataproc东西但是当从jdbc读取数据时它默认为1个分区,因为它是单线程的。我敢打赌你的 dataframe 是 1 个分区 = 1 个任务 = 1 个正在使用的核心,在这种情况下,只在一台机器上使用,因为没有任何东西被并行化。
  • 谢谢,我现在将 CPU 数量增加到 2 X 3 个工作程序,并将其重新分区为 10,但仍然没有改进。我忘记提到的另一件重要的事情是这个 json 文件是压缩文件 (.gz) 格式。不确定这是否会导致任何问题。
  • 是的,这是个问题... gzip 格式不可拆分... 可能与您的问题的解决方案在这里重复 => stackoverflow.com/questions/40492967/…

标签: json apache-spark pyspark google-cloud-storage google-cloud-dataproc


【解决方案1】:

最好的办法是对输入进行预处理。给定一个输入文件,spark.read.json(... 将创建一个任务来读取和解析 JSON 数据,因为 Spark 无法提前知道如何并行化它。如果您的数据是行分隔的 JSON 格式 (http://jsonlines.org/),最好的做法是预先将其拆分为可管理的块:

path = "gs://dataproc-bucket/json-files/data_5M.json"
# read monolithic JSON as text to avoid parsing, repartition and *then* parse JSON
mi = spark.read.json(spark.read.text(path).repartition(1000).rdd)
inf_rel = mi.select(mi.client_id,
                   mi.user_id,
                   mi.first_date,
                   F.hour(mi.first_date).alias('hour'),
                   mi.notes).rdd.map(foo).filter(lambda x: x)
inf_relevance = inf_rel.map(lambda l: Row(**dict(l))).toDF()
save_path = "gs://dataproc-bucket/json-files/output_5M.json"
inf_relevance.write.mode('append').json(save_path)
print("END!!")

您在此处的初始步骤 (spark.read.text(...) 仍将成为单个任务的瓶颈。如果您的数据不是行分隔的,或者(尤其是!)您预计您将需要多次使用此数据,您应该想办法在使用 Spark 之前将您的 5GB JSON 文件转换为 1000 个 5MB JSON 文件。

【讨论】:

  • 该文件是 gz 压缩文件,在生产中,我必须读取 64 个这样的 json 压缩文件,每个文件大小约为 5GB,你认为进一步切碎它们会有帮助吗?
  • 是的。更多文件 = 更多并行性 = 更少解析开销。
【解决方案2】:

.gz 文件不可拆分,因此它们由一个内核读取并放置在单个分区上。

参考Dealing with a large gzipped file in Spark

【讨论】:

    猜你喜欢
    • 2019-10-23
    • 2016-02-28
    • 1970-01-01
    • 2016-07-20
    • 1970-01-01
    • 2017-10-29
    • 2016-01-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多