惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Engineering at Meta
Engineering at Meta
Microsoft Azure Blog
Microsoft Azure Blog
I
InfoQ
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
人人都是产品经理
人人都是产品经理
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
T
Tailwind CSS Blog
MongoDB | Blog
MongoDB | Blog
Google DeepMind News
Google DeepMind News
WordPress大学
WordPress大学
量子位
美团技术团队
大猫的无限游戏
大猫的无限游戏
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
Last Week in AI
Last Week in AI
博客园 - 司徒正美
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
小众软件
小众软件
C
Check Point Blog
博客园 - 三生石上(FineUI控件)
N
Netflix TechBlog - Medium
Recent Announcements
Recent Announcements
有赞技术团队
有赞技术团队
月光博客
月光博客

膨胀自留地

使用CloudFlare标记管理功能优雅的使用谷歌分析 四天轻松版北京旅游攻略 Adsense税务居住地证明申请教程:如何申请中国税收居民身份证明以及Google Adsense更新新加坡税务信息教程 常见的 VPS 优化措施 HTML实现table表格移动端自动转为卡片样式 个人网站接入 Google 登录 如何从 GitHub 永久删除泄露的 .env 文件 Cloudflare设置www的域名跳转到不带www的域名 顶级iOS神器!全面解开App Store限制,AssPP Pro保姆级教程 BroadcastChannel将你的 Telegram Channel 转为微博客 Medama开源轻量级网站统计(分析)系统 白嫖一个网站监控面板:uptime-kuma 部署TraffMonetizer利用闲置VPS挂机赚钱 Counterscale-使用CloudflareWorkers部署的免费网站数据统计工具,界面类似umami 封装一个最简单的Axios,避免过度封装! 如何使用Cloudflare Email Routing创建域名邮箱免费发送邮件 使用Cron定时更新Hosts解决Gihub无法访问问题 容器化MariaDB/Mysql实现定时备份 JetBrains全家桶许可证服务器怎么找,JetBrains激活教程 图片加速接口:缓存图片,加速访问,解决防盗链 MySQL模糊查询再也不用like+%了,全文索引介绍及使用简介 为什么这么多人推荐你办理流量卡?流量卡怎么赚钱?怎么避坑? 借助Cloudflare Worker实现Google Analytics反代加速,规避广告屏蔽插件的拦截 etcd入门之在Centos上安装部署 Grafana在linux下安装和基本使用入门 Prometheus介绍、安装和基本功能快速入门 无需加速器,油猴插件解锁xbox云游戏 阿里云盘自动每日签到,无需部署,无需服务器 自定义VSCode终端主题样式 一行 CSS 代码实现响应式布局 – 使用 Grid 实现的响应式布局
python使用多线程threading模块长期循环运行内存泄漏问题解决
2022-04-26 · via 膨胀自留地

#编程技术 2022-04-26 15:45:00 | 全文 1818 字,阅读约需 4 分钟 | 加载中... 次浏览

👋 相关阅读


threading 是 python 中一个内置的多线程库,关于模块这里不做过多介绍,本文主要记录一下遇到的问题及解决方法。

一直以来都是用如下方法执行多线程任务,正常情况下是没问题。

import threading, time

## 限制线程的最大数量
threadmax = threading.BoundedSemaphore(16)
## 将锁内的代码串行化
lock = threading.Lock()
## 多线程 join 用
l = []


def do_something(i):
    '''
    ## 上锁的例子
    :param i:
    :return:
    '''
    with lock:  ## 上锁
        print(f"当前第 {i} 个任务正在执行")
    time.sleep(5)
    ## 释放信号量,可用信号量加一
    threadmax.release()


def do_something_unlock(i):
    '''
    ## 不上锁的例子
    :param i:
    :return:
    '''
    print(f"当前第 {i} 个任务正在执行")
    time.sleep(5)
    ## 释放信号量,可用信号量加一
    threadmax.release()


if __name__ == '__main__':
    for i in range(1000):
        ## 增加信号量,可用信号量减一
        threadmax.acquire()
        t = threading.Thread(target=do_something, args=(i,))
        t.start()
        l.append(t)
    for t in l:
        t.join()

    print('END')

但是工作中遇到一个需求,程序需要 24 小时运行,并且一直循环遍历某个列表并发去执行任务,所以就简单的在上边例子的基础上加了个 while 循环,如下:

import threading, time

## 限制线程的最大数量
threadmax = threading.BoundedSemaphore(16)
## 多线程 join 用
l = []


def do_something(i):
    '''
    ## 执行任务
    :param i:
    :return:
    '''
    print(f"当前第 {i} 个任务正在执行")
    time.sleep(5)
    ## 释放信号量,可用信号量加一
    threadmax.release()


if __name__ == '__main__':
	while 1:
		for i in range(1000):
			## 增加信号量,可用信号量减一
			threadmax.acquire()
			t = threading.Thread(target=do_something, args=(i,))
			t.start()
			l.append(t)
		for t in l:
			t.join()
    	print('end once')

运行几天之后发现问题,内存一直在增加,逐行测试之后,发现问题出在 L 这个变量上。

L 这个变量主要是用来存放需要 join 的任务的,在第一个 for 循环中,每个线程开始执行之后,都会存放到 L 这个列表里,用来支持第二个 for 循环中的 join 调用。

join 本质是用于堵塞当前主线程的类,其作用是阻止全部的线程执行完之前程序继续往下运行,直到被调用的线程全部执行完毕或者超时。

举个例子,do_something 这个函数需要执行一段时间,那么你多线程执行这个函数之后,如果没有 join,那么第一个 for 循环执行完之后就会接着往后执行,不会等待 do_something 这个函数执行完,也就是会直接执行执行 print('end once') 这行代码,代码示例如下:

import threading, time

## 限制线程的最大数量
threadmax = threading.BoundedSemaphore(16)
## 多线程 join 用
l = []


def do_something(i):
    '''
    ## 执行任务
    :param i:
    :return:
    '''
    print(f"当前第 {i} 个任务正在执行")
    time.sleep(2)
    print(f"当前第 {i} 个任务执行完成")
    ## 释放信号量,可用信号量加一
    threadmax.release()


if __name__ == '__main__':
    for i in range(10):
        ## 增加信号量,可用信号量减一
        threadmax.acquire()
        t = threading.Thread(target=do_something, args=(i,))
        t.start()
    print('end once')


## 控制台输出
当前第 0 个任务正在执行
当前第 1 个任务正在执行
当前第 2 个任务正在执行
当前第 3 个任务正在执行
当前第 4 个任务正在执行
当前第 5 个任务正在执行
当前第 6 个任务正在执行
当前第 7 个任务正在执行
当前第 8 个任务正在执行
当前第 9 个任务正在执行
end once

当前第 9 个任务执行完成
当前第 6 个任务执行完成当前第 7 个任务执行完成
当前第 5 个任务执行完成
当前第 4 个任务执行完成当前第 3 个任务执行完成当前第 2 个任务执行完成
当前第 8 个任务执行完成



当前第 0 个任务执行完成
当前第 1 个任务执行完成

Process finished with exit code 0

可以看到,print('end once') 这行代码并不是在最后执行的,如果想让 print('end once') 等到所有线程都运行完成后再执行,那就需要 join, 代码示例如下:

import threading, time

## 限制线程的最大数量
threadmax = threading.BoundedSemaphore(16)
## 多线程 join 用
l = []


def do_something(i):
    '''
    ## 执行任务
    :param i:
    :return:
    '''
    print(f"当前第 {i} 个任务正在执行")
    time.sleep(2)
    print(f"当前第 {i} 个任务执行完成")
    ## 释放信号量,可用信号量加一
    threadmax.release()


if __name__ == '__main__':
    for i in range(10):
        ## 增加信号量,可用信号量减一
        threadmax.acquire()
        t = threading.Thread(target=do_something, args=(i,))
        t.start()
        l.append(t)
    for t in l:
        t.join()
    print('end once')

## 控制台输出
当前第 0 个任务正在执行
当前第 1 个任务正在执行
当前第 2 个任务正在执行
当前第 3 个任务正在执行
当前第 4 个任务正在执行
当前第 5 个任务正在执行
当前第 6 个任务正在执行
当前第 7 个任务正在执行
当前第 8 个任务正在执行
当前第 9 个任务正在执行
当前第 7 个任务执行完成当前第 9 个任务执行完成当前第 2 个任务执行完成当前第 1 个任务执行完成
当前第 4 个任务执行完成当前第 6 个任务执行完成
当前第 3 个任务执行完成

当前第 5 个任务执行完成


当前第 0 个任务执行完成

当前第 8 个任务执行完成
end once

Process finished with exit code 0

这样就达到了我们的目的,这种写法一般情况下没有问题,但是如果在外边又套了一层 while,那就有问题了,因为 L 这个列表一直是在增加的,没有释放,所以导致内存一直增加。

找到问题根本就好说了,给他加一个释放

import threading, time

## 限制线程的最大数量
threadmax = threading.BoundedSemaphore(16)
## 多线程 join 用
l = []


def do_something(i):
    '''
    ## 执行任务
    :param i:
    :return:
    '''
    print(f"当前第 {i} 个任务正在执行")
    time.sleep(5)
    ## 释放信号量,可用信号量加一
    threadmax.release()


if __name__ == '__main__':
	while 1:
		for i in range(1000):
			## 增加信号量,可用信号量减一
			threadmax.acquire()
			t = threading.Thread(target=do_something, args=(i,))
			t.start()
			l.append(t)
		for t in l:
			t.join()
			l.remove(t)
    	print('end once')

加上 l.remove(t) 这一行,经测试内存稳定,如果不需要等待的话,也可以直接不 join。


×