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
验证步骤
- 先添加MongoDB分区配置,通过Spark UI的Jobs页面查看分区数是否合理(建议在200-400之间)。
- 调整动态分配和网络超时参数,观察Executor的扩容/缩容是否正常。
- 逐步提升Driver资源,监控Driver的CPU和内存使用率(通过Kubernetes Dashboard或Spark UI的Driver页面)。
内容的提问来源于stack exchange,提问作者P.Subedi
相关产品推荐
相关产品推荐

