同一Executor内任务间是否支持可读写的本地共享变量?
Great question! You're right that Accumulators and Broadcast Variables have their limitations for this specific use case—Accumulators are write-only from tasks (read-only on the driver) and Broadcasts are read-only, plus both involve network transfer.
Luckily, there are two solid, widely-used approaches to create executor-local, read-write shared variables that live within each executor's JVM, are accessible to all tasks running on that executor, and don't require network overhead. Let's break them down:
1. Executor-Level Singleton Objects (Scala/Java)
Since each Spark executor runs as a single JVM process, you can leverage language-level singletons to create a shared state that's initialized once per executor and accessible to all tasks (threads) on that executor.
Example (Scala):
// This singleton is initialized once per executor JVM object ExecutorSharedStore { // Use a thread-safe collection to avoid concurrency issues private val sharedData = new java.util.concurrent.ConcurrentHashMap[String, String]() def setValue(key: String, value: String): Unit = { sharedData.put(key, value) } def getValue(key: String): Option[String] = { Option(sharedData.get(key)) } } // Usage in your Spark job val rdd = sc.parallelize(1 to 100, 10) rdd.foreach { num => val key = s"task_$num" ExecutorSharedStore.setValue(key, s"processed_by_${Thread.currentThread().getName}") // All tasks on this executor can read the value println(s"Value for $key: ${ExecutorSharedStore.getValue(key)}") }
Key Notes:
- Always use thread-safe data structures (like
ConcurrentHashMapinstead of a regularMap) or add synchronization (e.g.,synchronizedblocks) since multiple tasks run concurrently on the same executor. - The singleton is isolated to each executor—changes made by tasks on one executor won't affect other executors, which aligns with your requirement for per-executor independent instances.
2. Custom Executor Plugins (Spark 3.0+)
For more control over initialization (e.g., loading local files, setting up connections to local services), you can use Spark's ExecutorPlugin API, which lets you run code when an executor starts up.
Example (Java):
import org.apache.spark.api.plugin.ExecutorPlugin; import org.apache.spark.SparkEnv; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; public class LocalExecutorStorePlugin implements ExecutorPlugin { private static ConcurrentHashMap<String, Object> executorStore; @Override public void init(SparkEnv env, Map<String, String> config) { // This runs once when the executor starts executorStore = new ConcurrentHashMap<>(); // Optional: Load initial data from local files or config executorStore.put("init_config", config.getOrDefault("custom.init.key", "default_value")); } // Static methods for tasks to access the store public static void put(String key, Object value) { executorStore.put(key, value); } public static Object get(String key) { return executorStore.get(key); } }
To Use This Plugin:
When submitting your Spark job, add the configuration to enable the plugin:
spark-submit \ --class com.yourpackage.YourJob \ --conf spark.executor.plugins=com.yourpackage.LocalExecutorStorePlugin \ --conf custom.init.key="my_initial_value" \ your-job.jar
Why This Works:
- The
initmethod runs exactly once per executor, so you can set up expensive resources (like local cache connections) without reinitializing them for every task. - Tasks can call the static
put/getmethods to interact with the shared store directly, no network transfer needed.
How This Compares to Accumulators/Broadcasts
- Accumulators: Only support append operations from tasks, and values are only readable on the driver—they can't be used for task-to-task sharing on the same executor.
- Broadcast Variables: Read-only, and are transferred from the driver to executors over the network. They don't support write operations from tasks.
Both of the approaches above solve your requirements perfectly: per-executor independent instances, full read-write access for all tasks on the executor, and zero network overhead.
内容的提问来源于stack exchange,提问作者nikniknik

