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

如何并行化分页HTTP请求的do-while循环以提升数据采集效率?

Parallelizing Pagination Requests in Java (Thread-Safe Implementation)

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> with ConcurrentHashMap<String, String> (JDK's built-in thread-safe map)
  • Replace LinkedListMultimap<String, Category> with ConcurrentHashMultimap<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.
  • ObjectMapper Safety: Your static ObjectMapper is thread-safe once configured, so it's safe to use across all threads.
  • getCategory Check: Make sure getCategory(posting) doesn't rely on shared mutable state—if it does, add synchronization or refactor it to be thread-safe.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:36:19