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

多线程数据采集与处理架构实现求助:求代码/伪代码示例

Great question! Let's walk through how to implement this collector-worker pattern properly, focusing on the synchronization mechanisms that make this work reliably.

Core Requirements Recap

We need:

  • Multiple collector threads that gather data and push it to a shared buffer
  • A single worker thread that pulls data from the shared buffer, processes it, and stores results in its own private buffer
  • The worker must only exit after all collectors have finished their work AND every piece of data in the shared buffer has been processed

Key Mechanisms to Use

To make this work without race conditions or incorrect termination, we'll rely on three core tools:

  1. Thread-safe blocking queue: Acts as the shared buffer, handling synchronization between collectors and the worker automatically.
  2. Thread synchronization primitives: To track when all collectors have finished (like join() in Python or CountDownLatch in Java).
  3. Completion signaling: To let the worker know no new data will be added, so it can finish processing remaining data and exit.

Python Implementation Example

Python's threading and queue modules make this straightforward. Here's a working example:

import threading
import queue
import random
import time

# Simulate data collection (replace with your actual logic)
def collect_data():
    time.sleep(random.uniform(0.1, 0.5))  # Simulate collection delay
    return random.randint(1, 100)

# Collector thread logic
def collector(source_buffer, collector_id):
    print(f"Collector {collector_id} starting up")
    for _ in range(5):  # Each collector gathers 5 data points
        data = collect_data()
        source_buffer.put(data)
        print(f"Collector {collector_id} added data: {data}")
    print(f"Collector {collector_id} finished collecting")

# Worker thread logic
def worker(source_buffer, private_buffer, collectors_done):
    print("Worker starting up")
    while True:
        try:
            # Wait for data, with a timeout to check completion status
            data = source_buffer.get(timeout=0.5)
            # Simulate processing (replace with your actual logic)
            processed_data = data * 2
            private_buffer.append(processed_data)
            print(f"Worker processed: {data} -> {processed_data}")
            source_buffer.task_done()
        except queue.Empty:
            # Check if all collectors are done AND no data is left to process
            if collectors_done.is_set() and source_buffer.empty():
                print("Worker: All data processed, exiting")
                break
    print(f"Worker finished. Results: {private_buffer}")

if __name__ == "__main__":
    # Shared buffer (thread-safe)
    shared_buffer = queue.Queue(maxsize=10)
    # Worker's private buffer (only accessed by worker, no sync needed)
    worker_results = []

    # Event to signal all collectors have finished
    collectors_completed = threading.Event()

    # Start collector threads
    collectors = []
    for i in range(3):  # 3 collector threads
        thread = threading.Thread(target=collector, args=(shared_buffer, i))
        collectors.append(thread)
        thread.start()

    # Start worker thread
    worker_thread = threading.Thread(target=worker, args=(shared_buffer, worker_results, collectors_completed))
    worker_thread.start()

    # Wait for all collectors to finish
    for thread in collectors:
        thread.join()
    print("All collectors have finished their work")
    # Signal the worker that no new data will come
    collectors_completed.set()

    # Wait for the worker to finish processing
    worker_thread.join()

Java Implementation Example

If you're using Java, we'll use LinkedBlockingQueue and CountDownLatch for synchronization:

import java.util.concurrent.*;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;

public class CollectorWorkerDemo {
    public static void main(String[] args) throws InterruptedException {
        BlockingQueue<Integer> sharedBuffer = new LinkedBlockingQueue<>(10);
        List<Integer> workerResults = new ArrayList<>(); // Private to worker
        int collectorCount = 3;
        CountDownLatch collectorLatch = new CountDownLatch(collectorCount);
        Random random = new Random();

        # Start collector threads
        for (int i = 0; i < collectorCount; i++) {
            int collectorId = i;
            new Thread(() -> {
                System.out.println("Collector " + collectorId + " starting");
                try {
                    for (int j = 0; j < 5; j++) {
                        // Simulate data collection
                        Thread.sleep(random.nextInt(400) + 100);
                        int data = random.nextInt(100) + 1;
                        sharedBuffer.put(data);
                        System.out.println("Collector " + collectorId + " added data: " + data);
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } finally {
                    collectorLatch.countDown();
                    System.out.println("Collector " + collectorId + " finished");
                }
            }).start();
        }

        // Start worker thread
        Thread worker = new Thread(() -> {
            System.out.println("Worker starting");
            while (true) {
                try {
                    // Poll for data with timeout to check completion
                    Integer data = sharedBuffer.poll(500, TimeUnit.MILLISECONDS);
                    if (data != null) {
                        int processed = data * 2;
                        workerResults.add(processed);
                        System.out.println("Worker processed: " + data + " -> " + processed);
                    } else {
                        // Check if all collectors are done and buffer is empty
                        if (collectorLatch.getCount() == 0 && sharedBuffer.isEmpty()) {
                            System.out.println("Worker: All data processed, exiting");
                            break;
                        }
                    }
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                    break;
                }
            }
            System.out.println("Worker finished. Results: " + workerResults);
        });
        worker.start();

        // Wait for all collectors to finish
        collectorLatch.await();
        System.out.println("All collectors have finished");

        // Wait for worker to complete processing
        worker.join();
    }
}

Why This Works

  • Thread-safe queues: The shared buffer handles all locking for us, so we don't have to worry about multiple collectors writing or the worker reading at the same time.
  • Completion signaling: The Event (Python) or CountDownLatch (Java) tells the worker when no more data will be added to the buffer.
  • Dual condition check: The worker only exits when both all collectors are done AND the buffer is empty—this ensures we don't leave unprocessed data behind, and we don't wait forever for new data that will never come.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:51:29