嵌套ThreadPoolExecutor执行异常:计数跳变而非逐步递增问题排查
问题:两层多线程执行时计数异常跳变
场景说明
我用concurrent.futures.ThreadPoolExecutor实现两层多线程做子域名枚举:
- 第一层:遍历约10000个域名,线程池最大工作线程数设为50
- 第二层:每个域名运行7款命令行工具(通过
subprocess执行),线程池最大工作线程数设为7
预期计数应该逐步递增,到50后等待线程释放再继续,但实际计数到50后几秒内直接跳至4000+,完全不符合预期。
核心代码片段
main.py(主程序)
import enumFunctions domains = get_domains() # 字符串列表,存储待枚举域名 count=0 with concurrent.futures.ThreadPoolExecutor(max_workers=50) as executor: for domain in domains: count+=1 print(count) executor.submit(enumFunctions.subdomainEnum, domain[0],client)
enumFunctions.py(工具执行模块)
function_and_parameters = [ (tool1,domain) # tool1为工具函数 (tool2,domain) # tool2为工具函数 ... ] with concurrent.futures.ThreadPoolExecutor(max_workers=7) as executor: for tool,param in function_and_parameters: executor.submit(tool, param)
可复现最小代码
main.py
import sys import enumFunctions import time import os current_dir = os.path.dirname(os.path.abspath(__file__)) import logging import concurrent.futures offset=15000 count=0 chunck=5000 max_count=9999999 domains=[x for x in range(0,100000)] while True: with concurrent.futures.ThreadPoolExecutor(max_workers=50) as executor: for domain in domains: count+=1 print(count) executor.submit(enumFunctions.subdomainEnum, domain) offset+=chunck if offset>=max_count: offset = 0
enumFunctions.py
import random import time import concurrent.futures def subdomainEnum(domain): function_and_parameters = [ (tool1,domain), (tool2,domain), (tool3,domain), (tool4,domain), (tool5,domain) ] with concurrent.futures.ThreadPoolExecutor(max_workers=7) as executor: for func,p1 in function_and_parameters: print(domain,func) executor.submit(func,p1) def tool1(domain): time.sleep(random.randint(4,30)) print("tool1",domain) return def tool2(domain): time.sleep(random.randint(4,30)) print("tool2",domain) return def tool3(domain): time.sleep(random.randint(4,30)) print("tool3",domain) return def tool4(domain): time.sleep(random.randint(4,30)) print("tool4",domain) return def tool5(domain): time.sleep(random.randint(4,30)) print("tool5",domain) return
问题原因
submit方法非阻塞:executor.submit()仅将任务放入线程池的任务队列,不会等待任务执行。主程序的for循环会快速遍历所有域名,一次性把所有任务都提交到队列,因此count会瞬间暴涨,而非逐步递增。- 线程池
max_workers的误解:max_workers控制的是同时运行的线程数量,不是限制任务提交的批次。任务提交过程是同步完成的,所有任务会被快速推入队列,线程池只是从队列中取任务执行,不影响count的递增速度。 - 外层
with块的作用局限:with ThreadPoolExecutor仅会在代码块结束时等待所有任务执行完毕,但任务提交的过程在for循环里已经同步完成,所以count会一次性跑完所有域名的计数。
解决方案
如果要实现“计数到50后等待线程释放再继续”的分批执行逻辑,需要手动控制任务提交的批次:
import enumFunctions import concurrent.futures domains = get_domains() count = 0 batch_size = 50 # 每批提交50个任务 # 按批次遍历域名列表 for i in range(0, len(domains), batch_size): batch_domains = domains[i:i+batch_size] with concurrent.futures.ThreadPoolExecutor(max_workers=batch_size) as executor: for domain in batch_domains: count +=1 print(count) executor.submit(enumFunctions.subdomainEnum, domain[0], client) # 当前批次所有任务执行完成后,再进入下一批次
这样每批仅提交50个域名任务,等待这批任务全部执行完成后再提交下一批,计数就会按预期逐步递增。
内容的提问来源于stack exchange,提问作者zifan yan
相关产品推荐
相关产品推荐

