【问题标题】:Dataproc Didn't Process Big Data in Parallel Using pysparkDataproc 未使用 pyspark 并行处理大数据
【发布时间】:2021-05-22 12:09:57
【问题描述】:

我在 GCP 中启动了一个 DataProc 集群,有一个主节点和 3 个工作节点。每个节点有 8 个 vCPU 和 30G 内存。

我开发了一个 pyspark 代码,它从 GCS 读取一个 csv 文件。 csv文件大小约30G。

df_raw = (
    spark
        .read
        .schema(schema)
        .option('header', 'true')
        .option('quote', '"')
        .option('multiline', 'true')
        .csv(infile)
)
df_raw = df_raw.repartition(20, "Product")
print(df_raw.rdd.getNumPartitions())

这是我将 pyspark 启动到 dataproc 中的方式:

gcloud dataproc jobs submit pyspark gs://<my-gcs-bucket>/<my-program>.py \
    --cluster=${CLUSTER} \
    --region=${REGION} \

我得到的分区号只有1。

我在此处附上了节点使用情况图片供您参考。

似乎它只使用了来自一个工作节点的一个 vCore。

如何使其与多个分区并行并使用所有节点和更多 vCore?

尝试重新分区到20,但仍然只使用了一个工作节点的一个vCore,如下:

Pyspark 默认分区为 200。所以我很惊讶地看到 dataproc 没有将所有可用资源用于此类任务。

【问题讨论】:

  • 您的帖子缺少实际在 DF 上运行的代码。如果这不是并行代码,例如Spark SQL、foreachPartition 等,那么您的代码将只能在 master 上运行,而不会在 executor 上运行。

标签: apache-spark pyspark dataproc


【解决方案1】:

这不是数据处理问题,而是纯 Spark/pyspark 问题。

为了并行化您的数据,需要将其拆分为多个分区 - 一个大于您拥有的执行器数量(工作核心总数)的数量。 (例如 ~ *2, ~ *3, ...)

有多种方法可以做到这一点,例如:

  1. 将数据拆分为文件或文件夹并并行化文件/文件夹列表并处理每个文件(或使用已经执行此操作的数据库并将此分区保持在 Spark 读取中)。

  2. 在获得 Spark DF 后重新分区数据,例如读取执行程序的数量并将它们乘以 N 并重新分区到这么多分区。当你这样做时,你必须选择能很好地划分你的数据的列,即分成许多部分,而不是仅仅分成几个部分,例如按天,按客户 ID,而不是按状态 ID。

df = df.repartition(num_partitions, 'partition_by_col1', 'partition_by_col2')

代码在主节点上运行,并行阶段分布在工作节点之间,例如

df = (
    df.withColumn(...).select(...)...
    .write(...)
)

由于 Spark 函数是惰性的,它们仅在您执行诸如 write 或 collect 之类的导致 DF 被评估的步骤时运行。

【讨论】:

  • 试过但没用。在原始帖子中添加了详细信息。谢谢。
【解决方案2】:

您可能想尝试通过 Dataproc 命令行的--properties 传递 Spark 配置来增加执行器的数量。所以像

gcloud dataproc jobs submit pyspark gs://<my-gcs-bucket>/<my-program>.py \
    --cluster=${CLUSTER} \
    --region=${REGION} \
    --properties=spark.executor.instances=5

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-10-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多