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

Python中实现任务并行替代串行:实时传递ID处理的方案咨询

实现Task1与Task2的并行协作(生产者-消费者模式)

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 used None here) 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.Queue instead of queue.Queue
  • Use multiprocessing.Process instead of threading.Thread

4. Error Handling Tips

  • Wrap your extraction/processing logic in try-except blocks 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:05:38