优化面向多用户的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:
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
CompletableFuturethe moment new matching data is saved.
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
findByLastAndByUserquery has a composite index onpreviousMessageId(I assumedataIdin your code is a typo) anduserIdto 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.
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
ConcurrentHashMapto track pending requests byuserId + previousMessageId—if a user sends the same request twice, return the existingCompletableFutureinstead of spinning up a new thread. - Timeout handling: Add a timeout to your
CompletableFutureto avoid threads hanging indefinitely (e.g., if a user disconnects without canceling the request).
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.
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); } }
- 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

