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

多线程下基于线程锁的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.Queue handles 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 on result_queue.get() to retrieve it.
  • Graceful Shutdown: We send a None task 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:55:17