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

Spark 2.2 Shell小作业Executor使用异常及数据Locality问题咨询

Spark 2.2: Troubleshooting Partition/Executor Utilization & Data Locality Questions

Hey there, let’s work through your Spark issues one by one—this stuff can feel confusing at first, but we’ll break it down clearly!


Issue 1: WordCount Uses Only 2 Partitions/Executors (vs. 3 Cores in sc.range Job)

First, let’s unpack why these two jobs behave so differently:

Why sc.range uses all 3 cores

The sc.range method creates an RDD that directly honors your cluster’s parallelism settings. In standalone mode, Spark defaults spark.default.parallelism to the total number of available cores across your workers. Since you have 3 cores, Spark splits the range into 3 partitions—so every core gets a task to run, utilizing all resources as expected.

Why your 60MB WordCount has 2 partitions

For text files, Spark relies on Hadoop’s FileInputFormat to calculate partition count, which depends on a few key factors:

  1. HDFS Block Size: If your HDFS block size is set to ~30MB, a 60MB file will split into 2 blocks, resulting in 2 partitions.
  2. Split Size Configurations: Check these parameters (set via spark-shell arguments or spark-defaults.conf):
    • spark.hadoop.mapreduce.input.fileinputformat.split.maxsize: The maximum size of a single split (default is 128MB in Hadoop—if you’re seeing 2 partitions, this value might be lowered to something like 32MB).
    • spark.hadoop.mapreduce.input.fileinputformat.split.minsize: The minimum split size; if this is smaller than the block size, it can force further splits of existing blocks.
  3. File Structure: If your "60MB file" is actually two 30MB files in the same directory, Spark will create one partition per file, leading to 2 total partitions.

Why only 2 Executors are used

Spark assigns one task per partition. With 2 partitions, you only need 2 tasks to process the data. If your Executors are configured with 1 core each (the default in standalone mode when you don’t set --executor-cores), Spark will spin up 2 Executors to handle the tasks—leaving your third core idle because there’s no task to assign to it.

Fixes to use all 3 cores

  • Manually set partition count: When loading the text file, explicitly specify the number of partitions:
    val textRDD = sc.textFile("/path/to/your/60mb/file", 3)
    
  • Adjust split size parameters: Launch spark-shell with a smaller max split size to force 3 partitions for your 60MB file:
    bin/spark-shell --master spark://localhost:7077 --executor-memory 6G --conf spark.hadoop.mapreduce.input.fileinputformat.split.maxsize=20971520
    
    (20971520 bytes = 20MB, so 60MB / 20MB = 3 splits)
  • Increase executor cores: If you set --executor-cores 3, a single Executor could handle all 3 tasks (assuming you have enough memory to support it).

Issue 2: Data Locality Questions (Common Scenarios & Fixes)

Since you didn’t specify exact locality pain points, I’ll cover the most frequent issues and fixes for Spark 2.2:

What are Spark’s Locality Levels?

Spark prioritizes tasks based on how close they are to the data (from best to worst performance):

  • PROCESS_LOCAL: Data is in the same JVM as the task (ideal, no network transfer)
  • NODE_LOCAL: Data is on the same node as the task (minimal overhead)
  • RACK_LOCAL: Data is on the same rack as the task (moderate network transfer)
  • ANY: Data is on a different rack entirely (high network overhead)

Why does Locality Degrade?

  • Insufficient resources: If there aren’t enough free cores/Executors on the node holding the data, Spark will move the task to another node instead of waiting.
  • Data skew: If most of your data is concentrated on a few nodes, those nodes get overloaded, forcing tasks to run on other nodes.
  • Short wait time: Spark has a default 3-second timeout for waiting for local resources—if this is too short, it quickly falls back to a lower locality level.

How to Optimize Locality

  • Adjust spark.locality.wait: Increase this value to give Spark more time to find local resources before falling back. For example:
    bin/spark-shell --master spark://localhost:7077 --executor-memory 6G --conf spark.locality.wait=10s
    
  • Match Executors to data nodes: Ensure you have enough Executors running on every node that holds your data.
  • Fix data skew: Use techniques like salting (adding a random prefix to keys) to distribute skewed data evenly across partitions.
  • Cache data locally: Use rdd.cache() or rdd.persist(StorageLevel.MEMORY_ONLY) to keep frequently accessed data on the node where it’s processed, enabling PROCESS_LOCAL locality for subsequent tasks.

内容的提问来源于stack exchange,提问作者mingzhao.pro

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 12:05:26