You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Dataproc Serverless Batch读取GCS Parquet文件过慢问题解决

问题:Dataproc Serverless Batch读取GCS中Parquet文件速度远慢于本地

将维基百科语料库的逆频率Parquet文件存储在Google Cloud Storage(GCS)中,加载到Dataproc Serverless(Batch)环境时遇到明显性能差异:用pyspark.read在Dataproc Batch中读取文件耗时20-30秒,而本地16GB内存、8核Intel CPU的MacBook上加载并持久化耗时不到10秒。

涉及的inverse_freq.parquet文件大小为148.8MB,存储桶采用标准存储类,使用Dataproc Batch运行时2.0版本。测试约50MB的小型Parquet文件时,Dataproc Batch的读取耗时仍为20-30秒,推测是Dataproc Batch配置或Spark初始化存在问题。


相关配置信息

1. 自定义Docker镜像配置

# Debian 11 is recommended.
FROM debian:11-slim

# Suppress interactive prompts
ENV DEBIAN_FRONTEND=noninteractive

# (Required) Install utilities required by Spark scripts.
RUN apt update && apt install -y procps tini libjemalloc2
# RUN apt-key adv --keyserver keyserver.ubuntu.com --recv-keys B8F25A8A73EACF41

# Enable jemalloc2 as default memory allocator
ENV LD_PRELOAD=/usr/lib/x86_64-linux-gnu/libjemalloc.so.2

# (Optional) Add extra jars.
ENV SPARK_EXTRA_JARS_DIR=/opt/spark/jars/
ENV SPARK_EXTRA_CLASSPATH='/opt/spark/jars/*'
RUN mkdir -p "${SPARK_EXTRA_JARS_DIR}"
#COPY spark-bigquery-with-dependencies_2.12-0.22.2.jar "${SPARK_EXTRA_JARS_DIR}"

# (Optional) Install and configure Miniconda3.
ENV CONDA_HOME=/opt/miniconda3
ENV PYSPARK_PYTHON=${CONDA_HOME}/bin/python
ENV PATH=${CONDA_HOME}/bin:${PATH}
COPY Miniconda3-py39_4.10.3-Linux-x86_64.sh .
RUN bash Miniconda3-py39_4.10.3-Linux-x86_64.sh -b -p /opt/miniconda3 \
  && ${CONDA_HOME}/bin/conda config --system --set always_yes True \
  && ${CONDA_HOME}/bin/conda config --system --set auto_update_conda False \
  && ${CONDA_HOME}/bin/conda config --system --prepend channels conda-forge \
  && ${CONDA_HOME}/bin/conda config --system --set channel_priority strict

# (Optional) Install Conda packages.
# Use mamba to install packages quickly.
RUN ${CONDA_HOME}/bin/conda install mamba -n base -c conda-forge \
    && ${CONDA_HOME}/bin/mamba install \
      conda \
      google-cloud-logging \
      python

ENV REQUIREMENTSPATH=/opt/requirements/requirements.txt
COPY requirements.txt "${REQUIREMENTSPATH}"
RUN pip install -r "${REQUIREMENTSPATH}"

ENV NLTKDATA_PATH=${CONDA_HOME}/nltk_data/corpora
RUN bash -c 'mkdir -p $NLTKDATA_PATH/{stopwords,wordnet}'
COPY nltk_data/stopwords ${NLTKDATA_PATH}/stopwords
COPY nltk_data/wordnet ${NLTKDATA_PATH}/wordnet

# (Optional) Add extra Python modules.
ENV PYTHONPATH=/opt/python/packages
RUN mkdir -p "${PYTHONPATH}"
RUN bash -c 'mkdir -p $PYTHONPATH/{utils,GCP}'
COPY utils "$PYTHONPATH/utils"
COPY GCP "$PYTHONPATH/GCP"

# (Required) Create the 'spark' group/user.
# The GID and UID must be 1099. Home directory is required.
RUN groupadd -g 1099 spark
RUN useradd -u 1099 -g 1099 -d /home/spark -m spark
USER spark

2. 提交Dataproc Batch任务的gcloud CLI命令

APP_NAME="context-graph"
BUCKET="context-graph"
IDF_PATH='idf_data/idf_data/inverse_freq.parquet'
DOC_PATH="articles/text.txt"

gcloud dataproc batches submit pyspark main.py \
    --version 2.0\
    --batch test \
    --container-image "custom_image:tag1" \
    --project project_id \
    --region us-central1 \
    --deps-bucket context_graph_deps \
    --service-account account@example.com \
    --subnet default \
    --properties spark.dynamicAllocation.initialExecutors=2,spark.dynamicAllocation.minExecutors=2,spark.executor.cores=4,spark.driver.cores=8,spark.driver.memory='16g',\
spark.executor.heartbeatInterval=200s,spark.network.timeout=250s\
    -- --app-name=${APP_NAME} --idf-uri=gs://${BUCKET}/${IDF_PATH} \
    --bucket-name=${BUCKET} --doc-path=${DOC_PATH}

3. 读取Parquet文件的main.py代码

import time

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

start = time.time()

df = (
    spark.read.option("inferSchema", "true")
    .option("header", "true")
    .parquet("gs://bucket/inverse_freq.parquet")
)
df.persist()

end = time.time()
print("loading time:", end - start)

4. 日志信息

Cloud Dataproc Batch日志中存在相关警告与错误。


解决方案

在创建SparkSession时添加master("local[*]")参数,修改后的代码如下:

spark = SparkSession.builder.master("local[*]").config(conf=conf).getOrCreate()

该参数不仅能解决Parquet文件读取缓慢的问题,还能改善GCS上pyspark.ml模型管道的加载速度,建议在涉及GCS读写的操作中添加此参数。


内容的提问来源于stack exchange,提问作者Sam

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 14:00:35