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

基于作业内存需求动态调整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.

Custom Memory-Aware Executor Implementation

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 ManagementFactory to 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 finally block 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 09:01:55