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

优化面向多用户的REST异步轮询数据操作

Hey there! Let's tackle optimizing this async polling endpoint for your authenticated users—right now that loop-based approach can get pretty inefficient as your user base scales up. Here are some practical, battle-tested improvements you can implement:

1. Ditch Polling for Event-Driven Architecture

The biggest issue with your current loop is that it wastes thread resources: every waiting user holds an async thread hostage while it sleeps and checks for data repeatedly. Instead, switch to an event-driven model where you push new data to waiting users as soon as it's available, rather than making them ask for it.

For REST scenarios, this can mean:

  • Using Server-Sent Events (SSE) to keep a long-lived connection open and stream data when it arrives (great for frontend clients).
  • Registering listeners for data creation events, and completing the CompletableFuture the moment new matching data is saved.
2. Optimize Cache & Database Efficiency

Your current code hits the database first—let's tweak that to reduce load:

  • Check cache first: If you're using a cache (like Redis or Caffeine), query it before hitting the database. Make sure to invalidate/update the cache immediately when new data is saved.
  • Add database indexes: Ensure your findByLastAndByUser query has a composite index on previousMessageId (I assume dataId in your code is a typo) and userId to speed up lookups.
  • Avoid redundant queries: If multiple users are polling for the same criteria, don't run duplicate DB checks—reuse the result across requests.
3. Control Concurrency & Prevent Resource Exhaustion

Uncontrolled async threads will crash your app as user count grows. Fix this with:

  • Custom async thread pool: Don't rely on Spring's default executor. Define a bounded pool to limit the number of concurrent polling threads.
  • Request deduplication: Use a ConcurrentHashMap to track pending requests by userId + previousMessageId—if a user sends the same request twice, return the existing CompletableFuture instead of spinning up a new thread.
  • Timeout handling: Add a timeout to your CompletableFuture to avoid threads hanging indefinitely (e.g., if a user disconnects without canceling the request).
4. Remove Blocking Operations

If your loop uses Thread.sleep(), that's a blocking call that ties up threads. Replace it with non-blocking delays if you must keep a fallback polling mechanism (though event-driven is better). For example, use ScheduledExecutorService to schedule a single check instead of looping with sleep.

Example Optimized Code

Here's how you might refactor your service to use event-driven logic + request deduplication + thread pool control:

First, define a custom async thread pool:

@Configuration
@EnableAsync
public class AsyncPollConfig {
    @Bean(name = "pollDataExecutor")
    public TaskExecutor pollDataTaskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(10); // Adjust based on your server capacity
        executor.setMaxPoolSize(50);
        executor.setQueueCapacity(100);
        executor.setThreadNamePrefix("DataPoll-");
        executor.initialize();
        return executor;
    }
}

Then, refactor your polling service to use event-driven completion and deduplication:

@Service
public class DataPollingService {
    private final DataRepository dataRepository;
    private final ConcurrentHashMap<String, CompletableFuture<List<Data>>> pendingRequests = new ConcurrentHashMap<>();

    public DataPollingService(DataRepository dataRepository) {
        this.dataRepository = dataRepository;
    }

    @Async("pollDataExecutor")
    public CompletableFuture<List<Data>> pollData(Long previousMessageId, Long userId) {
        String requestKey = userId + ":" + previousMessageId;

        // Return existing pending request if duplicate
        if (pendingRequests.containsKey(requestKey)) {
            return pendingRequests.get(requestKey);
        }

        CompletableFuture<List<Data>> resultFuture = new CompletableFuture<>();
        pendingRequests.put(requestKey, resultFuture);

        // Check for existing new data first
        List<Data> newData = dataRepository.findByLastAndByUser(previousMessageId, userId);
        if (!newData.isEmpty()) {
            completeAndCleanup(requestKey, resultFuture, newData);
            return resultFuture;
        }

        // Set timeout to avoid hanging threads
        resultFuture.orTimeout(30, TimeUnit.SECONDS)
                .exceptionally(ex -> {
                    pendingRequests.remove(requestKey);
                    return Collections.emptyList();
                });

        return resultFuture;
    }

    // Call this method when new data is saved for a user
    public void onNewDataSaved(Long userId, List<Data> newData) {
        // Find all pending requests for this user and complete them
        pendingRequests.entrySet().removeIf(entry -> {
            String key = entry.getKey();
            if (key.startsWith(userId + ":")) {
                Long reqPrevId = Long.parseLong(key.split(":")[1]);
                List<Data> filteredData = newData.stream()
                        .filter(data -> data.getId() > reqPrevId)
                        .collect(Collectors.toList());
                
                if (!filteredData.isEmpty()) {
                    entry.getValue().complete(filteredData);
                    return true;
                }
            }
            return false;
        });
    }

    private void completeAndCleanup(String requestKey, CompletableFuture<List<Data>> future, List<Data> data) {
        future.complete(data);
        pendingRequests.remove(requestKey);
    }
}
Key Improvements in This Code
  • No more wasteful looping: Threads only exist if there's actual work to do.
  • Duplicate requests are merged, reducing DB and thread load.
  • Thread pool is bounded, preventing resource exhaustion.
  • Timeouts ensure threads don't hang forever.
  • New data triggers immediate completion of waiting futures.

内容的提问来源于stack exchange,提问作者Denis Stephanov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:43:09