linux线程池分析

一. 线程池学习文件

pool_test/  -> 线程池函数接口实现源码,简单实例。

系统编程项目接口设计说明书.doc  -> 详细说明了线程池各个函数的头文件/原型/参数/返回值..。

线程池模型.jpg  -> 帮助大家理解线程池原理。

二. 学习线程池实现过程?

1. 什么是线程池?

线程池就是多个线程组合起来的一个集合,当有任务时,线程就会处理任务,当没有任务时,线程休息。

2. 分析线程池源码

thread_pool.c  -> 线程池函数接口源码

thread_pool.h  -> 函数接口声明/结构体声明/头文件..

===========================================================

thread_pool.h

#define MAX_WAITING_TASKS  1000

-> 最大的任务等待个数

#define MAX_ACTIVE_THREADS 20     -> 最大线程个数

0)任务节点结构体

struct task

{

void *(*do_task)(void *arg);  -> 任务函数

void *arg;

-> 任务函数的参数

struct task *next;         -> 指向下一个任务节点的指针

};

1)线程池模型

typedef struct thread_pool

{

pthread_mutex_t lock;

-> 互斥锁

pthread_cond_t  cond;

-> 条件变量

bool shutdown;

-> 线程池关闭标识符号  true->关闭   false->未关闭

struct task *task_list;

-> 任务队列的头文件

pthread_t *tids;

-> 存放线程TID号空间地址

unsigned max_waiting_tasks;

-> 最大的等待任务的个数

unsigned waiting_tasks;  -> 当前等待任务的个数

unsigned active_threads;

-> 当前线程池中线程的个数

}thread_pool;

2)初始化线程池函数模型

bool init_pool(thread_pool *pool, unsigned int threads_number);

3)线程处理函数

void *routine(void *arg)

4)添加任务函数

bool add_task(thread_pool *pool,void *(*do_task)(void *arg), void *arg)

===========================================================

thread_pool.c

1)初始化线程池函数源码

2)线程处理函数源码

3)添加任务函数源码

源码:

头文件:

#ifndef _THREAD_POOL_H_
#define _THREAD_POOL_H_

#include <stdio.h>
#include <stdbool.h>
#include <unistd.h>
#include <stdlib.h>
#include <string.h>
#include <strings.h>

#include <errno.h>
#include <pthread.h>

#define MAX_WAITING_TASKS    1000
#define MAX_ACTIVE_THREADS    20

struct task
{
    void *(*do_task)(void *arg);
    void *arg;

    struct task *next;
};

typedef struct thread_pool
{
    pthread_mutex_t lock;
    pthread_cond_t  cond;
    bool shutdown;
    struct task *task_list;
    pthread_t *tids;
    unsigned max_waiting_tasks;
    unsigned waiting_tasks;
    unsigned active_threads;
}thread_pool;

bool init_pool(thread_pool *pool, unsigned int threads_number);
bool add_task(thread_pool *pool, void *(*do_task)(void *arg), void *task);
int  add_thread(thread_pool *pool, unsigned int additional_threads_number);
int  remove_thread(thread_pool *pool, unsigned int removing_threads_number);
bool destroy_pool(thread_pool *pool);

void *routine(void *arg);

#endif

功能函数:

#include "thread_pool.h"

void handler(void *arg)
{
    printf("[%u] is ended.\n",
        (unsigned)pthread_self());

    //解锁!
    pthread_mutex_unlock((pthread_mutex_t *)arg);
}

void *routine(void *arg)
{
    //接住线程池的地址
    thread_pool *pool = (thread_pool *)arg;
    struct task *p;

    while(1)
    {
        //取消例程函数,将来线程上锁了,如果收到取消请求,那么先解锁,再退出
        pthread_cleanup_push(handler, (void *)&pool->lock);

        //任务队列是属于临界资源。
        //访问任务队列之前都必须上锁。
        pthread_mutex_lock(&pool->lock);

        //如果当前线程池未被关闭并且线程池中没有需要处理的任务时:
        while(pool->waiting_tasks == 0 && !pool->shutdown)
        {
            //那么就进入条件变量中等待!
            pthread_cond_wait(&pool->cond, &pool->lock);
        }

        //如果线程等待任务为0,并且线程池已经关闭了。
        if(pool->waiting_tasks == 0 && pool->shutdown == true)
        {
            //解锁
            pthread_mutex_unlock(&pool->lock);    

            //走人
            pthread_exit(NULL);
        }

        //有任务做,代表肯定不是空链表,拿任务p
        p = pool->task_list->next;
        pool->task_list->next = p->next;

        //当前等待的任务的个数-1
        pool->waiting_tasks--;

        //解锁
        pthread_mutex_unlock(&pool->lock);

        //删除线程取消例程函数
        pthread_cleanup_pop(0);

        //设置线程不可以响应取消。
        pthread_setcancelstate(PTHREAD_CANCEL_DISABLE, NULL);

        //执行任务节点中函数过程中,不希望被别人取消掉。
        (p->do_task)(p->arg);

        //设置为可以响应取消
        pthread_setcancelstate(PTHREAD_CANCEL_ENABLE, NULL);

        //释放任务节点p的内存空间
        free(p);
    }

    pthread_exit(NULL);
}

bool init_pool(thread_pool *pool, unsigned int threads_number)
{
    //1. 初始化互斥锁
    pthread_mutex_init(&pool->lock, NULL);

    //2. 初始化条件变量
    pthread_cond_init(&pool->cond, NULL);

    //3. 线程池关闭标志为未关闭
    pool->shutdown = false;

    //4. 为任务队列头节点申请空间
    pool->task_list = malloc(sizeof(struct task));

    //5. 为线程TID号申请空间
    pool->tids = malloc(sizeof(pthread_t) * MAX_ACTIVE_THREADS);

    //错误判断
    if(pool->task_list == NULL || pool->tids == NULL)
    {
        perror("allocate memory error");
        return false;
    }

    //为任务队列头节点的指针域赋值NULL
    pool->task_list->next = NULL;

    //设置最大等待任务个数为1000
    pool->max_waiting_tasks = MAX_WAITING_TASKS;

    //设置当前等待任务的个数为0
    pool->waiting_tasks = 0;

    //设置当前线程池线程的个数
    pool->active_threads = threads_number;

    int i;
    //创建线程池中的子线程
    for(i=0; i<pool->active_threads; i++)
    {
        if(pthread_create(&((pool->tids)[i]), NULL,routine, (void *)pool) != 0)
        {
            perror("create threads error");
            return false;
        }
    }

    //初始化成功
    return true;
}

bool add_task(thread_pool *pool,void *(*do_task)(void *arg), void *arg)
{
    //为新节点申请内存空间
    struct task *new_task = malloc(sizeof(struct task));
    if(new_task == NULL)
    {
        perror("allocate memory error");
        return false;
    }

    //为新节点的数据域赋值
    new_task->do_task = do_task;  //函数
    new_task->arg = arg; //函数的参数

    //为新节点的指针域赋值
    new_task->next = NULL;  

    //访问任务队列前,先上锁!
    pthread_mutex_lock(&pool->lock);

    //如果当前等待任务个数>=1000,则添加任务失败!
    if(pool->waiting_tasks >= MAX_WAITING_TASKS)
    {
        //解锁
        pthread_mutex_unlock(&pool->lock);

        //输出错误信息
        fprintf(stderr, "too many tasks.\n");

        //释放刚刚初始化过的新节点
        free(new_task);

        return false;
    }

    //寻找任务队列的最后一个节点
    struct task *tmp = pool->task_list;
    while(tmp->next != NULL)
        tmp = tmp->next;

    //tmp->next = NULL;

    //把新节点尾插进去任务队列中
    tmp->next = new_task;

    //当前最大的等待任务的个数+1
    pool->waiting_tasks++;

    //解锁
    pthread_mutex_unlock(&pool->lock);

    //随机唤醒条件变量中其中一个线程起来就可以了。
    pthread_cond_signal(&pool->cond);

    return true;
}

int add_thread(thread_pool *pool, unsigned additional_threads)
{
    //如果新增0个线程
    if(additional_threads == 0)
        return 0; //直接返回0

    //添加后线程总数
    unsigned total_threads =
            pool->active_threads + additional_threads;

    int i, actual_increment = 0;

    //创建线程
    for(i = pool->active_threads;
          i < total_threads && i < MAX_ACTIVE_THREADS;
          i++)
       {
        if(pthread_create(&((pool->tids)[i]),NULL, routine, (void *)pool) != 0)
        {
            perror("add threads error");
            if(actual_increment == 0)
                return -1;

            break;
        }
        actual_increment++;  //真正创建的线程个数
    }

    //当前活跃的线程数 = 原来活跃的线程数 + 新实际创建的线程数
    pool->active_threads += actual_increment;

    return actual_increment;
}

int remove_thread(thread_pool *pool, unsigned int removing_threads)
{
    //如果需要删除0条线程
    if(removing_threads == 0)
        return pool->active_threads; //当前线程池活跃的线程个数

    //剩余的线程数 = 当前活跃的线程数 - 需要删除的线程。
    int remaining_threads = pool->active_threads - removing_threads;

    //线程池中至少有1条线程
    remaining_threads = remaining_threads > 0 ? remaining_threads : 1;

    int i;
    for(i=pool->active_threads-1; i>remaining_threads-1; i--)
    {
        errno = pthread_cancel(pool->tids[i]);
        if(errno != 0)
            break;
    }

    //如果取消失败,则函数返回-1
    if(i == pool->active_threads-1)
        return -1;
    else
    {
        //计算当前剩余实际的个数
        pool->active_threads = i+1;
        return i+1; //返回当前线程剩余的个数
    }
}

bool destroy_pool(thread_pool *pool)
{
    pool->shutdown = true; //当前线程池标志是关闭状态
    pthread_cond_broadcast(&pool->cond);
    int i;
    for(i=0; i<pool->active_threads; i++)
    {
        errno = pthread_join(pool->tids[i], NULL);

        if(errno != 0)
        {
            printf("join tids[%d] error: %s\n",
                    i, strerror(errno));
        }

        else
            printf("[%u] is joined\n", (unsigned)pool->tids[i]);

    }

    free(pool->task_list);
    free(pool->tids);
    free(pool);

    return true;
}

主函数:

#include "thread_pool.h"

void *mytask(void *arg) //线程的任务
{
    int n = (int)arg;

    //工作任务:余数是多少,就睡多少秒,睡完,任务就算完成
    printf("[%u][%s] ==> job will be done in %d sec...\n",
        (unsigned)pthread_self(), __FUNCTION__, n);

    sleep(n);

    printf("[%u][%s] ==> job done!\n",
        (unsigned)pthread_self(), __FUNCTION__);

    return NULL;
}

void *count_time(void *arg)
{
    int i = 0;
    while(1)
    {
        sleep(1);
        printf("sec: %d\n", ++i);
    }
}

int main(void)
{
    // 本线程用来显示当前流逝的秒数
    // 跟程序逻辑无关
    pthread_t a;
    pthread_create(&a, NULL, count_time, NULL);

    // 1, initialize the pool
    thread_pool *pool = malloc(sizeof(thread_pool));
    init_pool(pool, 2);
    //2个线程都在条件变量中睡眠

    // 2, throw tasks
    printf("throwing 3 tasks...\n");
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));

    // 3, check active threads number
    printf("current thread number: %d\n",
            remove_thread(pool, 0));//2
    sleep(9);

    // 4, throw tasks
    printf("throwing another 6 tasks...\n");
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));
    add_task(pool, mytask, (void *)(rand()%10));

    // 5, add threads
    add_thread(pool, 2);

    sleep(5);

    // 6, remove threads
    printf("remove 3 threads from the pool, "
           "current thread number: %d\n",
            remove_thread(pool, 3));

    // 7, destroy the pool
    destroy_pool(pool);
    return 0;
}

原文地址:https://www.cnblogs.com/zjlbk/p/11359578.html

时间: 2024-10-12 14:47:21

linux线程池分析的相关文章

Android-Universal-Image-Loader 学习笔记(五)线程池分析

UniveralImageLoader中的线程池 一般情况网络访问就需要App创建一个线程来执行(不然可能出现很臭的ANR),但是这也导致了当网络访问比较多的情况下,线程的数目可能指数增多,虽然Android系统理论上说可以创建无数个线程,但是某一时间段,线程数的急剧增加可能导致系统OOM. 在UIL中引入了线程池这种技术来管理线程.合理利用线程池能够带来三个好处. 第一:降低资源消耗.通过重复利用已创建的线程降低线程创建和销毁造成的消耗. 第二:提高响应速度.当任务到达时,任务可以不需要等到线

java线程池分析和应用

比较 在前面的一些文章里,我们已经讨论了手工创建和管理线程.在实际应用中我们有的时候也会经常听到线程池这个概念.在这里,我们可以先针对手工创建管理线程和通过线程池来管理做一个比较.通常,我们如果手工创建线程,需要定义线程执行对象,它实现的接口.然后再创建一个线程对象,将我们定义好的对象执行部分装载到线程中.对于线程的创建.结束和结果的获取都需要我们来考虑.如果我们需要用到很多的线程时,对线程的管理就会变得比较困难.我们手工定义线程的方式在时间和空间效率方面会存在着一些不足.比如说我们定义好的线程

Java 线程池分析

在项目中经常会用到java线程池,但是别人问起线程池的原理,线程池的策略怎么实现的? 答得不太好,所以按照源码分析一番,首先看下最常用的线程池代码: public class ThreadPoolTest { private static Executor executor= Executors.newFixedThreadPool(10); //一般通过fixThreadPool起线程池 public static void main(String[] args){ for(int i=0;i

Linux线程池在服务器上简单应用

一.问题描述 现在以C/S架构为例,客户端向服务器端发送要查找的数字,服务器端启动线程中的线程进行相应的查询,将查询结果显示出来. 二.实现方案 1. 整个工程以client.server.lib组织,如下图所示: 2. 进入lib, socket.h.socket.c /** @file socket.h @brief Socket API header file TCP socket utility functions, it provides simple functions that h

Linux线程池在server上简单应用

一.问题描写叙述 如今以C/S架构为例.client向server端发送要查找的数字,server端启动线程中的线程进行对应的查询.将查询结果显示出来. 二.实现方案 1. 整个project以client.server.lib组织.例如以下图所看到的: watermark/2/text/aHR0cDovL2Jsb2cuY3Nkbi5uZXQvd2FuZ3poaWNoZW5nMTk4Mw==/font/5a6L5L2T/fontsize/400/fill/I0JBQkFCMA==/dissolv

linux线程池

typedef struct task_queue { pthread_mutex_t mutex; pthread_cond_t cond; /* when no task, the manager thread wait for ;when a new task come, signal. */ struct task_node *head; /* point to the task_link. */ int number; /* current number of task, includ

JDK之线程池分析

线程池组成 一个线程池包括以下四个基本组成部分: 1.线程池管理器(ThreadPool):用于创建并管理线程池,包括 创建线程池,销毁线程池,添加新任务: 2.工作线程(PoolWorker):线程池中线程,在没有任务时处于等待状态,可以循环的执行任务: 3.任务接口(Task):每个任务必须实现的接口,以供工作线程调度任务的执行,它主要规定了任务的入口,任务执行完后的收尾工作,任务的执行状态等: 4.任务队列(taskQueue):用于存放没有处理的任务.提供一种缓冲机制. 线程池种类 1.

Linux线程池的实现

线程池的实现 1:自定义封装的条件变量 1 //condition.h 2 #ifndef _CONDITION_H_ 3 #define _CONDITION_H_ 4 5 #include <pthread.h> 6 7 typedef struct condition 8 { 9 pthread_mutex_t pmutex; 10 pthread_cond_t pcond; 11 }condition_t; 12 13 int condition_init(condition_t *c

Linux C编程之二十二 Linux线程池实现

一.线程池实现原理 1. 管理者线程 (1)计算线程不够用 创建线程 (2) 空闲线程太多 a. 销毁 更新要销毁的线程个数 通过条件变量完成的 b. 如果空闲太多,任务不够 线程阻塞在该条件变量上 c. 发送信号 pthread_cond_signal 2. 线程池中的线程 (1)从任务队列中取数据 任务队列任务 执行任务 (2)销毁空闲的线程 让线程执行pthread_exit 阻塞空闲的线程收到信号: 解除阻塞          只有一个往下执行          在执行任务之前做了销毁操