Python中实现任务并行替代串行:实时传递ID处理的方案咨询
Absolutely! This is exactly the kind of problem the producer-consumer pattern was built to solve. You can run Task1 (the "producer" extracting IDs) and Task2 (the "consumer" processing IDs) in parallel, with Task1 feeding IDs to Task2 the moment each one is extracted—no need to wait for all IDs to be pulled first.
In Python, the easiest way to pull this off uses built-in tools for thread-safe communication and parallel execution. Here's a complete, practical implementation tailored to your use case:
Step-by-Step Solution
1. Full Example Code
import threading import queue import time # Create a thread-safe queue to pass IDs between tasks # Set a maxsize to prevent memory overflow if processing lags behind extraction id_queue = queue.Queue(maxsize=15) # Sentinel value to tell the consumer when there are no more IDs to process STOP_SIGNAL = None def task1_extract_ids(): """Producer: Extract IDs from your 120GB document set""" # Replace this loop with your actual ID extraction logic for id_num in range(1, 11): # Simulating 10 extracted IDs print(f"Task1: Extracted ID {id_num}") id_queue.put(id_num) time.sleep(0.4) # Simulate time taken to pull one ID from docs # Send stop signal once all IDs are extracted id_queue.put(STOP_SIGNAL) print("Task1: Finished extracting all IDs") def task2_process_id(): """Consumer: Process each ID immediately upon receipt""" while True: current_id = id_queue.get() # Exit loop if we get the stop signal if current_id == STOP_SIGNAL: id_queue.task_done() break # Replace this with your actual ID processing logic print(f"Task2: Processing ID {current_id}") time.sleep(0.8) # Simulate time taken to process one ID id_queue.task_done() # Mark this ID's processing as complete print("Task2: Finished processing all IDs") if __name__ == "__main__": # Initialize and start both threads producer_thread = threading.Thread(target=task1_extract_ids) consumer_thread = threading.Thread(target=task2_process_id) producer_thread.start() consumer_thread.start() # Wait for producer to finish extracting all IDs producer_thread.join() # Wait for every ID in the queue to be processed id_queue.join() # Wait for consumer to clean up and exit consumer_thread.join() print("All tasks completed successfully!")
2. Key Components Explained
queue.Queue: Handles thread-safe communication between Task1 and Task2—no need to worry about race conditions when passing IDs.- Sentinel Value: The
STOP_SIGNAL(we usedNonehere) tells Task2 when there are no more incoming IDs, so it can exit gracefully instead of looping forever. - Thread Management: Using
join()ensures we wait for both tasks to finish properly before exiting the main program.id_queue.join()guarantees every extracted ID gets processed.
3. For CPU-Intensive Processing
If Task2 is CPU-heavy (not just IO-bound like reading/writing data), swap threading with multiprocessing to bypass Python's GIL. The pattern stays nearly identical:
- Use
multiprocessing.Queueinstead ofqueue.Queue - Use
multiprocessing.Processinstead ofthreading.Thread
4. Error Handling Tips
- Wrap your extraction/processing logic in
try-exceptblocks to handle individual failures without crashing the entire workflow. - Add a separate queue or list to log failed IDs if you need to review them later.
This approach will overlap the time spent extracting and processing IDs, cutting down your total runtime significantly compared to serial execution.
内容的提问来源于stack exchange,提问作者User54211

