多线程下基于线程锁的FIFO队列函数执行及结果获取方案问询
How to Run Function Calls in FIFO Queue Across Threads (with Return Values)
Got it, let's tackle this problem step by step. You need to execute function calls in strict FIFO order even when submitted from different threads, ensure only one function runs at a time, and capture each call's return value. Here's a clean, Python-based solution that hits all your requirements:
Core Concepts We'll Use
queue.Queue: Python's built-in thread-safe FIFO queue to hold pending tasks.- A dedicated worker thread: This thread pulls tasks from the queue one at a time, guaranteeing sequential execution.
- Per-task result queues: Each submitted task includes a small queue to send the return value back to the caller thread.
Full Implementation Code
import queue import threading import random import time # Your existing function definitions def func1(arg1): # Simulate work to make execution order visible time.sleep(random.uniform(0.1, 0.5)) return f"func1 returned: {arg1}" def func2(arg1): time.sleep(random.uniform(0.1, 0.5)) return f"func2 returned: {arg1}" # Global thread-safe task queue task_queue = queue.Queue() def worker(): """Dedicated thread to process tasks in FIFO order""" while True: task = task_queue.get() # Exit signal: stop worker if task is None if task is None: task_queue.task_done() break func, args, result_queue = task try: # Execute function and send result back result = func(*args) result_queue.put(result) except Exception as e: # Optional: pass exceptions back to caller result_queue.put(f"Error in {func.__name__}: {str(e)}") finally: # Mark task as completed task_queue.task_done() def submit_task(func, *args): """Submit a function call to the queue and return a queue for results""" result_queue = queue.Queue(maxsize=1) task_queue.put((func, args, result_queue)) return result_queue if __name__ == "__main__": # Start the worker thread (daemon so it exits with main thread) worker_thread = threading.Thread(target=worker, daemon=True) worker_thread.start() def random_caller(thread_id): """Simulate a thread that randomly calls func1 or func2""" for _ in range(3): func = random.choice([func1, func2]) arg = f"from thread {thread_id}" print(f"Thread {thread_id} submitted {func.__name__} with arg: {arg}") # Submit task and wait for result result_q = submit_task(func, arg) result = result_q.get() print(f"Thread {thread_id} received result: {result}") # Random delay between submissions time.sleep(random.uniform(0.1, 0.3)) # Launch 3 caller threads caller_threads = [] for i in range(3): t = threading.Thread(target=random_caller, args=(i,)) caller_threads.append(t) t.start() # Wait for all callers to finish submitting tasks for t in caller_threads: t.join() # Wait for all tasks in the queue to be processed task_queue.join() # Send exit signal to worker task_queue.put(None) worker_thread.join() print("All tasks completed!")
Key Details Explained
- Thread Safety:
queue.Queuehandles all locking internally, so we don't have to worry about race conditions when submitting tasks from multiple threads. - Strict FIFO Order: The worker thread processes tasks in the exact order they're added—no two functions run simultaneously, which meets your requirement.
- Return Value Capture: Each task gets its own
result_queue—once the worker finishes executing the function, it drops the result into this queue, and the caller can block onresult_queue.get()to retrieve it. - Graceful Shutdown: We send a
Nonetask to the worker to signal it should exit after all pending tasks are done.
Alternative Using concurrent.futures (Simpler)
If you prefer a higher-level approach, ThreadPoolExecutor with max_workers=1 automatically enforces sequential FIFO execution, and Future objects handle result retrieval cleanly:
from concurrent.futures import ThreadPoolExecutor import random import time def func1(arg1): time.sleep(random.uniform(0.1, 0.5)) return f"func1 returned: {arg1}" def func2(arg1): time.sleep(random.uniform(0.1, 0.5)) return f"func2 returned: {arg1}" if __name__ == "__main__": # 1 worker ensures sequential FIFO execution with ThreadPoolExecutor(max_workers=1) as executor: def random_caller(thread_id): futures = [] for _ in range(3): func = random.choice([func1, func2]) arg = f"from thread {thread_id}" print(f"Thread {thread_id} submitted {func.__name__} with arg: {arg}") # Submit task and store future for result future = executor.submit(func, arg) futures.append(future) time.sleep(random.uniform(0.1, 0.3)) # Retrieve results for this thread's tasks for future in futures: result = future.result() print(f"Thread {thread_id} received result: {result}") caller_threads = [] for i in range(3): t = threading.Thread(target=random_caller, args=(i,)) caller_threads.append(t) t.start() for t in caller_threads: t.join() print("All tasks completed!")
This version is more concise—ThreadPoolExecutor manages the queue and worker thread for you, so you don't have to write low-level queue handling code.
内容的提问来源于stack exchange,提问作者majidse
相关产品推荐
相关产品推荐

