使用Asyncio与aiohttp批量发送请求遇任务未完成问题求助
问题:asyncio+aiohttp实现大请求量脚本任务无法完成,并发限制失效
我要写一个脚本向网站发送10万到100万量级的请求,尝试用asyncio和aiohttp实现,但任务创建后无法完成,而且我对这两个库不太熟悉。
代码逻辑:先询问并发数和代理类型,调用start_workers创建启动任务,用asyncio.Semaphore限制同时并发10000个请求,但这个限制没生效;check函数负责发请求、处理响应和更新统计数据,这部分也不正常;另外写了console函数每2秒监控进度,count_requests_per_minute统计每分钟请求量。
附上完整代码:
import threading import os import time import random from queue import Queue from tkinter.filedialog import askopenfilename from tkinter import Tk from colorama import Fore from urllib.parse import quote import asyncio import aiohttp stats_lock = asyncio.Lock() stats = { 'valid': 0, 'invalid': 0, 'twofa': 0, 'error': 0, 'total_checked': 0, 'cpm': 0 } async def check(data, proxy, stats_lock, stats, session): payload = { 'data':data } headers2 = { "User-Agent": "Mozilla/5.0 (Linux; Android 6.0; Nexus 5 Build/MRA58N) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Mobile Safari/537.36", } async with session.post("https://example_website.com", headers=headers2, data=payload, proxy=proxy) as response: response_text = await response.text() if "true" in response_text: stats['valid'] +=1 elif "false" in response_text: stats['invalid']+=1 async def handler(data, proxy, proxytype, stats_lock, stats, session): async with session: proxy_url = f"http://{proxy}" await check(data, proxy_url, stats_lock, stats, session) async def start_workers(threads, data_queue, proxies_list, proxies_input): sem = asyncio.Semaphore(threads) console_thread = threading.Thread(target=console, args=(len(data_queue),)) console_thread.start() cpm_thread = threading.Thread(target=count_requests_per_minute) cpm_thread.start() tasks = [] async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(), trust_env=True) as session: for data in data_queue: task = asyncio.ensure_future(handler(data, random.choice(proxies_list), proxies_input, stats_lock, stats, session)) tasks.append(task) await asyncio.sleep(0) async with sem: await task async def main(): data_queue = [] root = Tk() root.withdraw() threads = int(input(f"{Fore.RED}[{Fore.WHITE}?{Fore.RED}]{Fore.WHITE} - {Fore.RED}How many threads {Fore.WHITE}:{Fore.RED}")) proxies_input = input(f"\n{Fore.RED}[{Fore.WHITE}!{Fore.RED}]{Fore.WHITE} > {Fore.RED}Proxies type {Fore.WHITE}| {Fore.RED}HTTP{Fore.WHITE}/{Fore.RED}SOCKS4{Fore.WHITE}/{Fore.RED}SOCKS5 {Fore.WHITE}:{Fore.RED} ") combo_file = askopenfilename(title="Data File", parent=root) with open(combo_file, "r", encoding='utf-8') as combofile: data_queue.extend(combofile.read().splitlines()) proxy_file = askopenfilename(title="Proxy File") with open(proxy_file, "r") as proxyfile: proxies_list = [line.strip() for line in proxyfile] await start_workers(threads, data_queue, proxies_list, proxies_input) def count_requests_per_minute(): while True: time.sleep(1) stats['cpm'] = stats['total_checked'] * 60 stats['total_checked'] = 0 def console(combo): print(f"\n{Fore.RED}[{Fore.WHITE}!{Fore.RED}]{Fore.WHITE} - {Fore.RED}Please wait while the console is loading.") time.sleep(10) os.system("cls") while True: os.system("cls") print(f""" {Fore.RED}[{Fore.WHITE}Valid{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['valid']} {Fore.RED}[{Fore.WHITE}Invalid{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['invalid']} {Fore.RED}[{Fore.WHITE}Errors{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['error']} {Fore.RED}[{Fore.WHITE}Checked{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['valid']+stats['invalid']}/{combo} {Fore.RED}[{Fore.WHITE}CPM{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['cpm']} """) time.sleep(2) if __name__ == "__main__": os.system("cls") try: loop = asyncio.get_event_loop() loop.run_until_complete(main()) except RuntimeError as e: if str(e) != "Event loop is closed": raise e except Exception: pass
问题分析与修复方案
核心问题梳理
- Semaphore并发限制完全失效:当前代码在循环里逐个创建任务后立刻
await task,变成串行执行,完全没用到并发;而且Semaphore的使用方式错误,没有在任务内部获取信号量。 - 统计数据线程不安全:
stats字典被asyncio协程和普通线程(console、count_requests_per_minute)同时修改,会出现数据竞争,导致统计错误。 - ClientSession滥用:
handler里又嵌套了async with session,但外层已经创建了全局的ClientSession,重复关闭会导致请求失败。 - CPM统计逻辑错误:每分钟请求量计算方式错误,而且
total_checked没有在协程里正确累加。 - 异常处理缺失:网络请求过程中没有捕获异常,导致单个请求失败就可能终止整个协程。
修复后的完整代码
import threading import os import time import random from tkinter.filedialog import askopenfilename from tkinter import Tk from colorama import Fore import asyncio import aiohttp from threading import Lock # 用线程锁保证stats的线程安全,因为普通线程和协程都会访问 stats_lock = Lock() stats = { 'valid': 0, 'invalid': 0, 'twofa': 0, 'error': 0, 'total_checked': 0, 'cpm': 0 } async def check(data, proxy, session): payload = {'data': data} headers = { "User-Agent": "Mozilla/5.0 (Linux; Android 6.0; Nexus 5 Build/MRA58N) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Mobile Safari/537.36", } try: async with session.post( "https://example_website.com", headers=headers, data=payload, proxy=proxy, timeout=aiohttp.ClientTimeout(total=10) ) as response: response_text = await response.text() with stats_lock: if "true" in response_text: stats['valid'] += 1 elif "false" in response_text: stats['invalid'] += 1 stats['total_checked'] += 1 except Exception as e: with stats_lock: stats['error'] += 1 stats['total_checked'] += 1 async def worker(data, proxy, sem, session): # 任务内部获取信号量,实现并发限制 async with sem: proxy_url = f"http://{proxy}" await check(data, proxy_url, session) async def start_workers(threads, data_queue, proxies_list): # 启动监控线程 console_thread = threading.Thread(target=console, args=(len(data_queue),), daemon=True) console_thread.start() cpm_thread = threading.Thread(target=count_requests_per_minute, daemon=True) cpm_thread.start() # 设置TCP连接池参数,避免Too many open files错误 connector = aiohttp.TCPConnector( limit_per_host=threads//10, limit=threads, ttl_dns_cache=300 ) async with aiohttp.ClientSession(connector=connector) as session: sem = asyncio.Semaphore(threads) tasks = [] # 批量创建任务,统一等待完成 for data in data_queue: proxy = random.choice(proxies_list) task = asyncio.create_task(worker(data, proxy, sem, session)) tasks.append(task) # 等待所有任务完成 await asyncio.gather(*tasks) async def main(): data_queue = [] root = Tk() root.withdraw() threads = int(input(f"{Fore.RED}[{Fore.WHITE}?{Fore.RED}]{Fore.WHITE} - {Fore.RED}并发数设置 {Fore.WHITE}:{Fore.RED}")) proxies_type = input(f"\n{Fore.RED}[{Fore.WHITE}!{Fore.RED}]{Fore.WHITE} > {Fore.RED}代理类型 {Fore.WHITE}| {Fore.RED}HTTP{Fore.WHITE}/{Fore.RED}SOCKS4{Fore.WHITE}/{Fore.RED}SOCKS5 {Fore.WHITE}:{Fore.RED} ") combo_file = askopenfilename(title="选择数据文件", parent=root) with open(combo_file, "r", encoding='utf-8') as combofile: data_queue.extend(combofile.read().splitlines()) proxy_file = askopenfilename(title="选择代理文件") with open(proxy_file, "r") as proxyfile: proxies_list = [line.strip() for line in proxyfile if line.strip()] await start_workers(threads, data_queue, proxies_list) def count_requests_per_minute(): while True: time.sleep(60) # 每分钟统计一次 with stats_lock: stats['cpm'] = stats['total_checked'] stats['total_checked'] = 0 def console(total): print(f"\n{Fore.RED}[{Fore.WHITE}!{Fore.RED}]{Fore.WHITE} - {Fore.RED}控制台加载中,请稍候。") time.sleep(2) os.system("cls" if os.name == "nt" else "clear") while True: os.system("cls" if os.name == "nt" else "clear") with stats_lock: checked = stats['valid'] + stats['invalid'] + stats['error'] print(f""" {Fore.RED}[{Fore.WHITE}有效{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['valid']} {Fore.RED}[{Fore.WHITE}无效{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['invalid']} {Fore.RED}[{Fore.WHITE}错误{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['error']} {Fore.RED}[{Fore.WHITE}已处理{Fore.RED}]{Fore.WHITE} | {Fore.RED}{checked}/{total} {Fore.RED}[{Fore.WHITE}每分钟请求量{Fore.RED}]{Fore.WHITE} | {Fore.RED}{stats['cpm']} """) time.sleep(2) if __name__ == "__main__": os.system("cls" if os.name == "nt" else "clear") try: asyncio.run(main()) except Exception as e: print(f"{Fore.RED}程序异常: {e}")
关键修复说明
- 并发限制修复:把Semaphore放到
worker函数内部,用async with sem包裹任务逻辑,同时批量创建所有任务后用asyncio.gather统一等待,真正实现并发。 - 线程安全修复:将
asyncio.Lock换成普通的threading.Lock,因为stats被普通线程和协程共享,asyncio.Lock只能在协程中使用,普通线程无法获取。 - ClientSession优化:只创建一个全局的ClientSession,设置合理的TCP连接池参数,避免文件句柄耗尽。
- CPM统计修复:改为每60秒统计一次过去一分钟的请求量,逻辑更准确。
- 异常处理:在
check函数中捕获所有请求异常,更新错误统计,避免单个请求失败导致整个协程崩溃。 - 守护线程设置:给监控线程加上
daemon=True,主程序退出时自动终止监控线程。
内容的提问来源于stack exchange,提问作者b4ddev
相关产品推荐
相关产品推荐

