如何并行化分页HTTP请求的do-while循环以提升数据采集效率?
Absolutely, you can parallelize this pagination logic to drastically cut down runtime when dealing with hundreds of pages. Since your HttpClient is already a thread-safe singleton, we just need to focus on thread-safe data collection and efficient parallel task execution. Here's a tailored solution for your code:
1. Swap Out Non-Thread-Safe Collections
Your current HashMap and LinkedListMultimap aren't designed for concurrent writes—multiple threads updating them will cause race conditions and data corruption. Replace them with thread-safe alternatives:
- Replace
HashMap<String, String>withConcurrentHashMap<String, String>(JDK's built-in thread-safe map) - Replace
LinkedListMultimap<String, Category>withConcurrentHashMultimap<String, Category>(from Guava, optimized for concurrent multimap operations)
Update your class fields like this:
private final Map<String, String> processToTaskIdHolder = new ConcurrentHashMap<>(); private final Multimap<String, Category> itemsByCategory = ConcurrentHashMultimap.create();
2. Refactor the Collection Workflow
We'll split the process into two clear phases to keep things manageable:
- Phase 1: Fetch the first page synchronously to get the total number of pages
- Phase 2: Parallelize requests for remaining pages using a controlled thread pool
Implementation with ExecutorService (Controlled Parallelism)
This approach lets you explicitly set how many concurrent requests to run (e.g., 5 at a time), which is great for respecting server rate limits:
private void collect() { String endpoint = "url_endpoint"; int totalPages = fetchFirstPageAndGetTotalPages(endpoint); if (totalPages <= 1) { return; // No more pages to process } // Create a fixed thread pool (adjust size based on your server's tolerance) try (ExecutorService executor = Executors.newFixedThreadPool(5)) { List<Callable<Void>> pageTasks = new ArrayList<>(); // Create tasks for pages 2 through totalPages for (int pageNumber = 2; pageNumber <= totalPages; pageNumber++) { pageTasks.add(() -> { fetchAndProcessSinglePage(endpoint, pageNumber); return null; }); } // Run all tasks and wait for them to finish executor.invokeAll(pageTasks); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Pagination collection was interrupted", e); } } // Helper to fetch the first page, process its data, and return total pages private int fetchFirstPageAndGetTotalPages(String endpoint) { HttpEntity<String> requestEntity = new HttpEntity<>(getBody(1), getHeader()); ResponseEntity<String> responseEntity = HttpClient.getInstance().getClient() .exchange(URI.create(endpoint), HttpMethod.POST, requestEntity, String.class); String jsonInput = responseEntity.getBody(); Stuff response = objectMapper.readValue(jsonInput, Stuff.class); // Process the first page's data (reuses your original logic) processPageData(response, jsonInput); return (int) response.getPaginationResponse().getNumberOfPages(); } // Reusable method to fetch and process a single page private void fetchAndProcessSinglePage(String endpoint, int pageNumber) { try { HttpEntity<String> requestEntity = new HttpEntity<>(getBody(pageNumber), getHeader()); ResponseEntity<String> responseEntity = HttpClient.getInstance().getClient() .exchange(URI.create(endpoint), HttpMethod.POST, requestEntity, String.class); String jsonInput = responseEntity.getBody(); Stuff response = objectMapper.readValue(jsonInput, Stuff.class); processPageData(response, jsonInput); } catch (IOException e) { // Add retry logic here if your API has transient errors throw new RuntimeException("Failed to process page " + pageNumber, e); } } // Extract data processing logic into a reusable, thread-safe method private void processPageData(Stuff response, String jsonInput) { List<Postings> postings = response.getPostings(); for (Postings posting : postings) { if (posting.getClientIds().isEmpty()) { continue; } List<String> lastParent = JsonPath.read(jsonInput, lastParentIdJsonPath); String clientId = posting.getClientIds().get(0).getId(); Category category = getCategory(posting); // Thread-safe writes thanks to concurrent collections itemsByCategory.put(clientId, category); processToTaskIdHolder.put(clientId, lastParent.get(0)); } }
Alternative: Modern Async with CompletableFuture
If you prefer a more concise, non-blocking approach, CompletableFuture works perfectly:
private void collect() { String endpoint = "url_endpoint"; int totalPages = fetchFirstPageAndGetTotalPages(endpoint); if (totalPages <= 1) { return; } // Use a fixed thread pool for controlled parallelism Executor executor = Executors.newFixedThreadPool(5); // Generate and run all async page tasks List<CompletableFuture<Void>> futures = IntStream.rangeClosed(2, totalPages) .mapToObj(pageNumber -> CompletableFuture.runAsync( () -> fetchAndProcessSinglePage(endpoint, pageNumber), executor )) .collect(Collectors.toList()); // Wait for all tasks to complete CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); }
3. Critical Notes for Safety & Performance
- Rate Limiting: Don't set the thread pool size too high (stick to 5-10 unless you confirm the server allows more). Too many concurrent requests can trigger rate limits or cause timeouts.
- Error Handling: Add retry logic (e.g., using Spring's
RetryTemplate) for transient errors like connection timeouts—this ensures you don't lose data from occasional failed requests. - Resource Cleanup: When using
ExecutorService, always use try-with-resources to ensure the pool shuts down properly after processing. ObjectMapperSafety: Your staticObjectMapperis thread-safe once configured, so it's safe to use across all threads.getCategoryCheck: Make suregetCategory(posting)doesn't rely on shared mutable state—if it does, add synchronization or refactor it to be thread-safe.
内容的提问来源于stack exchange,提问作者user1950349

