多线程数据采集与处理架构实现求助:求代码/伪代码示例
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:
- Thread-safe blocking queue: Acts as the shared buffer, handling synchronization between collectors and the worker automatically.
- Thread synchronization primitives: To track when all collectors have finished (like
join()in Python orCountDownLatchin Java). - 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) orCountDownLatch(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
相关产品推荐
相关产品推荐

