一. 引言

   如web服务器、数据库服务器等应用都面向处理来自某些远程来源的大量短小的任务。构建服务器应用程序的一个过于简单的模型是:每当一个请求到达就创建一个新的服务对象,然后在新的服务 对象中为请求服务。但当有大量请求并发访问时,服务器不断的创建和销毁对象的开销很大。所以提高服务器效率的一个手段就是尽可能减少创建和销毁对象的次 数,特别是一些很耗资源的对象创建和销毁,这样就引入了“池”的概念,“池”的概念使得人们可以定制一定量的资源,然后对这些资源进行复用,而不是频繁的 创建和销毁。

二. 基本概念

   线程池是预先创建线程的一种技术。线程池在还没有任务到来之前,创建一定数量的线程,放入空闲队列中。这些线程都是处于睡眠状态,即均为启动,不消耗 CPU,而只是占用较小的内存空间。当请求到来之后,缓冲池给这次请求分配一个空闲线程,把请求传入此线程中运行,进行处理。当预先创建的线程都处于运行 状态,即预制线程不够,线程池可以自由创建一定数量的新线程,用于处理更多的请求。当系统比较闲的时候,也可以通过移除一部分一直处于停用状态的线程。

线程池的设计思路:

   一个典型的线程池,应该包括如下几个部分:

  1、线程池管理器(ThreadPool),用于启动、停用,管理线程池
  2、工作线程(WorkThread),线程池中的线程
  3、请求接口(WorkRequest),创建请求对象,以供工作线程调度任务的执行
  4、请求队列(RequestQueue),用于存放和提取请求
  5、结果队列(ResultQueue),用于存储请求执行后返回的结果

  线程池管理器,通过添加请求的方法(putRequest)向请求队列(RequestQueue)添加请求,这些请求事先需要实现请求接口,即传递工作 函数、参数、结果处理函数、以及异常处理函数。之后初始化一定数量的工作线程,这些线程通过轮询的方式不断查看请求队列(RequestQueue),只 要有请求存在,则会提取出请求,进行执行。然后,线程池管理器调用方法(poll)查看结果队列(resultQueue)是否有值,如果有值,则取出, 调用结果处理函数执行。通过以上讲述,不难发现,这个系统的核心资源在于请求队列和结果队列,工作线程通过轮询requestQueue获得人物,主线程 通过查看结果队列,获得执行结果。因此,对这个队列的设计,要实现线程同步,以及一定阻塞和超时机制的设计,以防止因为不断轮询而导致的过多cpu开销。python的Queue,就是很好的实现了对线程同步机制。

三.  实现方式: 

  在python3.2版本之前,可以通过Queue + threading的方式。3.2版本以后,可以通过 concurrent.futures.ThreadPoolThreadPoolExecutor实现。 二者区别:

工具 核心定位 设计目的
ThreadPoolExecutor 高层线程池管理工具 简化线程池的创建、任务提交和结果获取,自动管理线程生命周期,适合批量提交任务,无需关心线程创建细节,自动分配任务到线程。
threading,Queue 线程安全的队列数据结构

实现多线程间的任务传递和通信,需手动管理线程创建与协作。

适合复杂的任务分发逻辑

1. concurrent.futures.ThreadPoolThreadPoolExecutor实现案例:

from concurrent.futures import ThreadPoolExecutor
import time

def task(n):
    time.sleep(1)
    return n * 2

# 线程池自动管理线程
with ThreadPoolExecutor(max_workers=3) as executor:
    # 提交任务
    results = executor.map(task, [1, 2, 3, 4, 5])
    
    # 获取结果
    for res in results:
        print(res)  # 输出:2,4,6,8,10

 适合的场景: 

  • 任务逻辑简单,只需批量执行函数
  • 不需要复杂的线程间协作
  • 希望代码简洁,减少样板代码



2. threading+ Queue 实现案例:

import threading
from queue import Queue
import time

def worker(queue):
    """工作线程:从队列取任务并执行"""
    while True:
        n = queue.get()  # 阻塞等待任务
        if n is None:  # 退出信号
            break
        time.sleep(1)
        print(n * 2)
        queue.task_done()  # 标记任务完成

# 创建队列
queue = Queue()

# 创建并启动3个工作线程
threads = []
for _ in range(3):
    t = threading.Thread(target=worker, args=(queue,))
    t.start()
    threads.append(t)

# 提交任务到队列
for n in [1, 2, 3, 4, 5]:
    queue.put(n)

queue.join()  # 等待所有任务完成

# 发送退出信号
for _ in range(3):
    queue.put(None)

# 等待线程结束
for t in threads:
    t.join()

 适合的场景: 

  • 需要自定义任务分发逻辑(如动态添加任务、优先级处理)
  • 实现生产者 - 消费者模型(如爬虫中的 URL 队列)
  • 需精细控制线程行为(如动态调整工作线程数量)

在复杂场景中,两者也可以结合使用。例如,用队列缓存任务,线程池从队列中取任务执行:

from concurrent.futures import ThreadPoolExecutor
from queue import Queue
import time

def process_task(queue):
    while True:
        n = queue.get()
        if n is None:
            break
        time.sleep(1)
        print(f"处理结果: {n * 2}")
        queue.task_done()

# 创建任务队列
task_queue = Queue()

# 线程池从队列取任务
with ThreadPoolExecutor(max_workers=3) as executor:
    # 启动工作线程
    for _ in range(3):
        executor.submit(process_task, task_queue)
    
    # 生产者添加任务
    for n in [1, 2, 3, 4, 5]:
        task_queue.put(n)
    
    task_queue.join()  # 等待所有任务完成
    
    # 发送退出信号
    for _ in range(3):
        task_queue.put(None)

建议: 简单场景优先选 ThreadPoolExecutor,复杂场景(如动态任务、优先级)选 threading.Queue,或两者结合使用。

 ---------------------------------------------------------------------------------------------------------------------

                         深耕运维行业多年,擅长运维体系建设,方案落地。欢迎交流!

                                                     “V-x”: ywjw996

                                                     《 运维经纬 》
 

更多推荐