Django+Celery调用Telethon报错RuntimeError:线程无当前事件循环
问题描述
我正在开发一个Telegram群组用户数据解析Web应用,通过POST表单接收群组用户名,解析该群组用户数据并保存到文件。将解析逻辑迁移到Celery后,在Django视图中调用该任务函数时,出现错误RuntimeError: There is no current event loop in thread 'Thread-1'。单独运行解析代码时一切正常。
tasks.py
from telethon.sync import TelegramClient import csv from celery import shared_task @shared_task def main(url): api_id = ******** api_hash = '********' phone = '******' client = TelegramClient(phone, api_id, api_hash) client.connect() target_group = client.get_entity(url) print('Fetching Members...') all_participants = [] all_participants = client.get_participants(target_group, aggressive=True) print('Saving In file...') with open("members.csv","w",encoding='UTF-8') as f: writer = csv.writer(f,delimiter=",",lineterminator="\n") writer.writerow(['username','user id', 'access hash','name','group', 'group id']) for user in all_participants: if user.username: username= user.username else: username= "" if user.first_name: first_name= user.first_name else: first_name= "" if user.last_name: last_name= user.last_name else: last_name= "" name= (first_name + ' ' + last_name).strip() writer.writerow([username ,user.id, user.access_hash, name, target_group.title, target_group.id]) print('Members scraped successfully.')
views.py
from django.shortcuts import render from django.views.decorators.csrf import csrf_exempt from .tasks import main @csrf_exempt def join_the_group(request): if request.method == 'POST': url = request.POST['group-name'] main(url) return render(request, 'ParseUserData/main_form.html') else: return render(request, 'ParseUserData/main_form.html')
错误堆栈信息
Traceback (most recent call last): File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/django/core/handlers/exception.py", line 55, in inner response = get_response(request) File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/django/core/handlers/base.py", line 197, in _get_response response = wrapped_callback(request, *callback_args, **callback_kwargs) File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/django/views/decorators/csrf.py", line 56, in wrapper_view return view_func(*args, **kwargs) File "/home/dev/Desktop/DjangoParser/ParserJSON/ParseUserData/views.py", line 10, in join_the_group main() File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/celery/local.py", line 182, in __call__ return self._get_current_object()(*a, **kw) File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/celery/app/task.py", line 411, in __call__ return self.run(*args, **kwargs) File "/home/dev/Desktop/DjangoParser/ParserJSON/ParseUserData/tasks.py", line 12, in main client = TelegramClient(phone, api_id, api_hash) File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/telethon/client/telegrambaseclient.py", line 335, in __init__ if not callable(getattr(self.loop, 'sock_connect', None)): File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/telethon/client/telegrambaseclient.py", line 483, in loop return helpers.get_running_loop() File "/home/dev/Desktop/DjangoParser/venv/lib/python3.8/site-packages/telethon/helpers.py", line 433, in get_running_loop return asyncio.get_event_loop_policy().get_event_loop() File "/usr/lib/python3.8/asyncio/events.py", line 639, in get_event_loop raise RuntimeError('There is no current event loop in thread %r.' RuntimeError: There is no current event loop in thread 'Thread-1'. [20/Aug/2023 01:32:00] "POST /user-data/collection-name HTTP/1.1" 500 104535
解决方案
问题根源
Celery工作线程默认不会初始化asyncio事件循环,而Telethon的sync模块只是对异步API的同步封装,底层依然依赖asyncio事件循环。单独运行代码时,主线程会自动创建事件循环,但Celery的工作线程没有这个机制,因此触发报错。
修复方案
方案1:手动创建并设置事件循环
修改tasks.py,在初始化TelegramClient前手动创建asyncio事件循环:
from telethon.sync import TelegramClient import csv from celery import shared_task import asyncio @shared_task def main(url): # 手动创建并绑定事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) api_id = ******** api_hash = '********' phone = '******' client = TelegramClient(phone, api_id, api_hash) client.connect() target_group = client.get_entity(url) print('Fetching Members...') all_participants = [] all_participants = client.get_participants(target_group, aggressive=True) print('Saving In file...') with open("members.csv","w",encoding='UTF-8') as f: writer = csv.writer(f,delimiter=",",lineterminator="\n") writer.writerow(['username','user id', 'access hash','name','group', 'group id']) for user in all_participants: username = user.username or "" first_name = user.first_name or "" last_name = user.last_name or "" name = (first_name + ' ' + last_name).strip() writer.writerow([username ,user.id, user.access_hash, name, target_group.title, target_group.id]) print('Members scraped successfully.') # 关闭事件循环 loop.close()
方案2:改用Telethon异步API重构任务
使用Telethon原生异步API,并用asyncio.run在Celery任务中执行:
from telethon import TelegramClient import csv from celery import shared_task import asyncio @shared_task def main(url): async def scrape_group(): api_id = ******** api_hash = '********' phone = '******' async with TelegramClient(phone, api_id, api_hash) as client: target_group = await client.get_entity(url) print('Fetching Members...') all_participants = [] async for user in client.iter_participants(target_group, aggressive=True): all_participants.append(user) print('Saving In file...') with open("members.csv","w",encoding='UTF-8') as f: writer = csv.writer(f,delimiter=",",lineterminator="\n") writer.writerow(['username','user id', 'access hash','name','group', 'group id']) for user in all_participants: username = user.username or "" first_name = user.first_name or "" last_name = user.last_name or "" name = (first_name + ' ' + last_name).strip() writer.writerow([username ,user.id, user.access_hash, name, target_group.title, target_group.id]) print('Members scraped successfully.') asyncio.run(scrape_group())
额外注意事项
调用Celery任务时必须使用delay()或apply_async()方法,否则会在当前线程同步执行,失去Celery异步处理的意义。修改views.py中的调用代码:
# 替换原来的 main(url) main.delay(url)
内容的提问来源于stack exchange,提问作者Grinya
相关产品推荐
相关产品推荐

