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

嵌套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

问题原因

  1. submit方法非阻塞:executor.submit()仅将任务放入线程池的任务队列,不会等待任务执行。主程序的for循环会快速遍历所有域名,一次性把所有任务都提交到队列,因此count会瞬间暴涨,而非逐步递增。
  2. 线程池max_workers的误解:max_workers控制的是同时运行的线程数量,不是限制任务提交的批次。任务提交过程是同步完成的,所有任务会被快速推入队列,线程池只是从队列中取任务执行,不影响count的递增速度。
  3. 外层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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 21:14:48