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

EKS上EMR Spark处理超大规模数据集时任务阻塞问题排查

问题解答

针对你在EKS上的EMR Spark处理7000万+行MongoDB数据时出现的Driver与Executor调度阻塞问题,结合你提供的配置,以下是缺失的关键配置及优化方案:

你的现有配置

Spark Submit 配置

spark_submit_cmd = [
        "--deploy-mode",
        "cluster",
        "--packages",
        "org.mongodb.spark:mongo-spark-connector_2.12:3.0.1,org.apache.hadoop:hadoop-aws:3.3.4,",  # noqa
        "--conf",
        f"spark.kubernetes.container.image={cfg('spark.docker.image')}",
        "--conf",
        "spark.dynamicAllocation.enabled=true",
        "--conf",
        "spark.dynamicAllocation.minExecutors=2",
        "--conf",
        "spark.dynamicAllocation.maxExecutors=5",
        "--conf",
        "spark.executor.memory=2g",
        "--conf",
        "spark.executor.cores=1",
        "--conf",
        f"spark.kubernetes.driver.podTemplateFile=s3://grepsr-aws-emr{'-stg' if os.getenv('APP_ENV')=='stg' else ''}/pod-template/spark-template.yaml",  # noqa
        "--conf",
        f"spark.kubernetes.executor.podTemplateFile=s3://grepsr-aws-emr{'-stg' if os.getenv('APP_ENV')=='stg' else ''}/pod-template/spark-template.yaml",  # noqa
        # "--conf",
        # "spark.kubernetes.executor.annotation.scheduler.alpha.kubernetes.io/tolerations=[{'key': 'role', 'operator': 'Equal', 'value': 'spark', 'effect': 'NoSchedule'}]" # noqa
        "--conf",
        "spark.driver.memory=8g",
        "--conf",
        "spark.driver.cores=1",
        "--conf",
        "spark.network.timeout=60s",
        "--conf",
        "spark.executor.heartbeatInterval=5s",
        "--conf",
        "spark.default.parallelism=100",
        "--conf",
        "spark.sql.shuffle.partitions=1000",
    ]


    for item in env_vars_items.keys():
        spark_submit_cmd.append("--conf")
        spark_submit_cmd.append(f"spark.kubernetes.driverEnv.CFG_{item.upper()}={env_vars_items[item]}")
        spark_submit_cmd.append("--conf")
        spark_submit_cmd.append(f"spark.executorEnv.CFG_{item.upper()}={env_vars_items[item]}")
        spark_submit_cmd.append("--conf")
        spark_submit_cmd.append(f"spark.kubernetes.executorEnv.CFG_{item.upper()}={env_vars_items[item]}")

SparkSession 配置

SparkSession.builder.appName(self.app_name)
            # .config("spark.master", os.getenv("CFG_SPARK_MASTER_HOST"))
            .config(
                "spark.mongodb.input.uri",
                os.getenv("CFG_SPARK_MONGO_URI"),
            )
            .config("spark.mongodb.input.database", os.getenv("CFG_MONGO_DB"))
            .config("spark.mongodb.input.collection", self.mongo_collection)
            .config(
                "spark.jars.packages",
                "org.mongodb.spark:mongo-spark-connector_2.12:3.0.1,org.apache.hadoop:hadoop-aws:3.3.4,",
            )
            .config("spark.hadoop.fs.s3a.access.key", os.getenv("CFG_AWS_ACCESS_KEY"))
            .config("spark.hadoop.fs.s3a.secret.key", os.getenv("CFG_AWS_SECRET_KEY"))
            .config("spark.hadoop.fs.s3a.endpoint", os.getenv("CFG_AWS_S3_ENDPOINT"))
            .config("spark.hadoop.fs.s3a.region", "eu-central-1")
            .config(
                "spark.hadoop.fs.s3a.aws.credentials.provider", "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
            )
            .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
            .config("spark.hadoop.fs.s3a.path.style.access", "true")
            # .config(
            #     "spark.hadoop.fs.s3a.block.size",
            #     os.getenv("CFG_AWS_S3_BLOCK_SIZE"),
            # )
            # .config(
            #     "spark.hadoop.fs.s3a.multipart.size",
            #     os.getenv("CFG_AWS_S3_MULTIPART_SIZE"),
            # )
            # .config("spark.hadoop.fs.s3a.committer.name", "/commits")
            # .config("spark.hadoop.fs.s3a.fast.upload", "true")
            # .config("spark.sql.sources.commitProtocolClass", "org.apache.hadoop.fs.s3a.commit.staging.S3ACommitterFactory") # noqa
            # .config("spark.hadoop.fs.s3a.committer.staging.conflict-mode", "replace")
            .getOrCreate()

关键缺失配置与优化方案

1. MongoDB读取分区与批次优化

7000万行数据的读取分区是核心瓶颈,现有配置未指定MongoDB连接器的分区策略,导致分区不均或数量不足:
在SparkSession中添加以下配置:

.config("spark.mongodb.input.partitioner", "MongoPaginateBySizePartitioner")
.config("spark.mongodb.input.partitionerOptions.partitionSizeMB", "64")
.config("spark.mongodb.input.batchSize", "10000")
  • MongoPaginateBySizePartitioner:按数据大小分区,比默认的_id范围分区更均衡,适合非分片集合。
  • partitionSizeMB:每个分区控制在64MB,匹配你的Executor 2G内存规格,避免OOM。
  • batchSize:单次读取1万条数据,减少MongoDB的请求次数,提升读取效率。

2. 动态分配完整配置

现有仅开启了动态分配,但缺少调度超时参数,导致Executor无法及时扩容/缩容:
在spark_submit_cmd中添加:

"--conf", "spark.dynamicAllocation.schedulerBacklogTimeout=10s",
"--conf", "spark.dynamicAllocation.executorIdleTimeout=60s",
"--conf", "spark.dynamicAllocation.cachedExecutorIdleTimeout=300s",
  • schedulerBacklogTimeout:任务积压10秒后触发Executor扩容。
  • executorIdleTimeout:Executor空闲60秒后释放,避免资源浪费。
  • cachedExecutorIdleTimeout:带缓存数据的Executor保留5分钟,减少重复读取。

3. Driver资源扩容

现有Driver仅1核,处理7000万行的任务调度和元数据会成为瓶颈:
修改spark_submit_cmd中的Driver配置:

"--conf", "spark.driver.cores=2",
"--conf", "spark.driver.maxResultSize=4g",
  • 提升Driver核心数到2,增强调度能力。
  • maxResultSize限制单个任务结果大小,避免Driver内存溢出导致调度中断。

4. 网络超时调整

现有60秒网络超时过短,大分区任务处理时Executor可能无法及时发送心跳:
修改spark_submit_cmd中的网络配置:

"--conf", "spark.network.timeout=300s",
"--conf", "spark.executor.heartbeatInterval=10s",

延长网络超时到5分钟,同时调整心跳间隔为10秒,平衡监控频率和负载。

5. 并行度与分区数匹配

现有spark.default.parallelism=100与spark.sql.shuffle.partitions=1000不匹配,导致shuffle阶段资源浪费:
修改spark_submit_cmd中的配置:

"--conf", "spark.default.parallelism=200",
"--conf", "spark.sql.shuffle.partitions=200",

并行度设置为Executor最大数核数的倍数(51*4=200),确保每个Executor有足够的任务处理。

6. Kubernetes Pod模板补充

检查你的podTemplate.yaml,确保Executor的资源请求/限制与Spark配置一致,添加存活探针避免被误杀:

containers:
- name: spark-kubernetes-executor
  resources:
    requests:
      memory: "2Gi"
      cpu: "1"
    limits:
      memory: "2Gi"
      cpu: "1"
  livenessProbe:
    exec:
      command:
        - /bin/sh
        - -c
        - "spark-submit --version"
    initialDelaySeconds: 30
    periodSeconds: 60

验证步骤

  1. 先添加MongoDB分区配置,通过Spark UI的Jobs页面查看分区数是否合理(建议在200-400之间)。
  2. 调整动态分配和网络超时参数,观察Executor的扩容/缩容是否正常。
  3. 逐步提升Driver资源,监控Driver的CPU和内存使用率(通过Kubernetes Dashboard或Spark UI的Driver页面)。

内容的提问来源于stack exchange,提问作者P.Subedi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 07:05:58