基于作业内存需求动态调整Java Executor执行的实现咨询
Great question—this is a common gap in the standard Java Executor framework, since it only handles thread count limits, not resource-based constraints like memory. Let's walk through a practical, maintainable implementation that solves this.
The core idea is to build a wrapper around a standard ExecutorService that tracks memory usage of running jobs and enforces admission control based on your light/heavy job classification. Here's a step-by-step breakdown:
1. Define Job Memory Classification
First, formalize your light/heavy job types with predefined memory cost estimates (adjust these values to match your actual workload):
public enum JobMemoryType { LIGHT(1), // Represents 1 unit of memory HEAVY(5); // Represents 5 units of memory private final int memoryCost; JobMemoryType(int memoryCost) { this.memoryCost = memoryCost; } public int getMemoryCost() { return memoryCost; } }
Then create a wrapper to bind tasks to their memory type:
public class MemoryAwareTask<T> implements Callable<T> { private final Callable<T> delegate; private final JobMemoryType memoryType; // For Callable tasks public MemoryAwareTask(Callable<T> delegate, JobMemoryType memoryType) { this.delegate = delegate; this.memoryType = memoryType; } // For Runnable tasks public MemoryAwareTask(Runnable delegate, JobMemoryType memoryType) { this(() -> { delegate.run(); return null; }, memoryType); } public JobMemoryType getMemoryType() { return memoryType; } @Override public T call() throws Exception { return delegate.call(); } }
2. Build the Memory-Tracking Executor Wrapper
This class wraps a standard ExecutorService (like ThreadPoolExecutor) and handles memory quota checks before submitting tasks, plus cleanup after tasks finish:
import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; public class MemoryAwareExecutorService implements ExecutorService { private final ExecutorService delegateExecutor; private final int totalMemoryQuota; private final AtomicInteger currentMemoryUsage = new AtomicInteger(0); public MemoryAwareExecutorService(ExecutorService delegateExecutor, int totalMemoryQuota) { this.delegateExecutor = delegateExecutor; this.totalMemoryQuota = totalMemoryQuota; } // Try to reserve memory for a new task private boolean tryAcquireMemory(int requiredMemory) { while (true) { int current = currentMemoryUsage.get(); int newTotal = current + requiredMemory; if (newTotal > totalMemoryQuota) { return false; // Not enough memory to accept the task } // Atomically update usage to avoid race conditions if (currentMemoryUsage.compareAndSet(current, newTotal)) { return true; } } } // Release memory when a task finishes private void releaseMemory(int releasedMemory) { currentMemoryUsage.addAndGet(-releasedMemory); } // Core submit method for MemoryAwareTasks @Override public <T> Future<T> submit(Callable<T> task) { if (!(task instanceof MemoryAwareTask)) { throw new IllegalArgumentException("Only MemoryAwareTask instances are supported"); } MemoryAwareTask<T> memoryTask = (MemoryAwareTask<T>) task; int requiredMemory = memoryTask.getMemoryType().getMemoryCost(); if (!tryAcquireMemory(requiredMemory)) { throw new RejectedExecutionException( String.format("Insufficient memory quota for %s task (current usage: %d/%d)", memoryTask.getMemoryType(), currentMemoryUsage.get(), totalMemoryQuota) ); } // Wrap task to release memory on completion (even if it throws an exception) return delegateExecutor.submit(() -> { try { return memoryTask.call(); } finally { releaseMemory(requiredMemory); } }); } // Convenience methods for direct task submission public Future<?> submit(Runnable task, JobMemoryType memoryType) { return submit(new MemoryAwareTask<>(task, memoryType)); } public <T> Future<T> submit(Callable<T> task, JobMemoryType memoryType) { return submit(new MemoryAwareTask<>(task, memoryType)); } // Delegate all other ExecutorService methods to the underlying executor @Override public void shutdown() { delegateExecutor.shutdown(); } @Override public List<Runnable> shutdownNow() { return delegateExecutor.shutdownNow(); } @Override public boolean isShutdown() { return delegateExecutor.isShutdown(); } @Override public boolean isTerminated() { return delegateExecutor.isTerminated(); } @Override public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException { return delegateExecutor.awaitTermination(timeout, unit); } // Disable raw submit/execute methods to enforce memory type usage @Override public Future<?> submit(Runnable task) { throw new UnsupportedOperationException("Use submit(Runnable, JobMemoryType) instead"); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException { throw new UnsupportedOperationException("Use submit with MemoryAwareTask instead"); } @Override public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException { throw new UnsupportedOperationException("Use submit with MemoryAwareTask instead"); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException { throw new UnsupportedOperationException("Use submit with MemoryAwareTask instead"); } @Override public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException { throw new UnsupportedOperationException("Use submit with MemoryAwareTask instead"); } @Override public void execute(Runnable command) { throw new UnsupportedOperationException("Use submit(Runnable, JobMemoryType) instead"); } }
3. Optional: Flexible Waiting with Semaphores
If you want tasks to wait for memory instead of being immediately rejected, replace the AtomicInteger memory tracker with a Semaphore:
private final Semaphore memorySemaphore; public MemoryAwareExecutorService(ExecutorService delegateExecutor, int totalMemoryQuota) { this.delegateExecutor = delegateExecutor; this.memorySemaphore = new Semaphore(totalMemoryQuota); } private boolean tryAcquireMemory(int requiredMemory) throws InterruptedException { // Use tryAcquire with timeout for non-blocking waits, or acquire() for blocking return memorySemaphore.tryAcquire(requiredMemory, 1, TimeUnit.SECONDS); } private void releaseMemory(int releasedMemory) { memorySemaphore.release(releasedMemory); }
4. Usage Example
public class MemoryExecutorDemo { public static void main(String[] args) { // Create a base thread pool (adjust core/max threads to match your workload) ExecutorService basePool = new ThreadPoolExecutor(2, 5, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(10)); // Initialize memory-aware executor with a total quota of 10 units MemoryAwareExecutorService memoryExecutor = new MemoryAwareExecutorService(basePool, 10); // Submit 8 light tasks (total memory: 8*1=8) for (int i = 0; i < 8; i++) { int taskId = i; memoryExecutor.submit(() -> { System.out.printf("Running light task %d%n", taskId); Thread.sleep(1000); }, JobMemoryType.LIGHT); } // Try to submit a heavy task (needs 5 units; 8+5=13 > 10, so it's rejected) try { memoryExecutor.submit(() -> { System.out.println("Running heavy task"); Thread.sleep(2000); }, JobMemoryType.HEAVY); } catch (RejectedExecutionException e) { System.out.println("Rejected: " + e.getMessage()); } // Wait for some light tasks to finish, then submit a heavy task try { Thread.sleep(3000); memoryExecutor.submit(() -> { System.out.println("Running heavy task after light tasks completed"); Thread.sleep(2000); }, JobMemoryType.HEAVY); } catch (Exception e) { e.printStackTrace(); } // Cleanup memoryExecutor.shutdown(); } }
Key Notes
- Estimated vs Actual Memory: This implementation uses fixed memory estimates. For more accuracy, you could integrate
ManagementFactoryto track actual thread memory usage, but this adds overhead—balance precision with performance. - Thread Pool Alignment: Ensure your base thread pool's size doesn't create CPU contention even when memory is available.
- Exception Safety: The
finallyblock guarantees memory is released even if tasks throw exceptions. - Dynamic Quotas: You could extend this to adjust the total memory quota based on real-time JVM available memory.
内容的提问来源于stack exchange,提问作者Crass Developer

