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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 16:47:16