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
相关产品推荐
相关产品推荐

