Discord.py异步代码在Airflow任务中无法关闭连接的问题
解决Airflow中Discord.py无法关闭连接、无法进入下一个任务的问题
问题描述
尝试将DataFrame发布至Discord频道时,遇到Discord.py无法关闭连接、无法进入下一个任务的问题。已尝试使用事件循环及asyncio.run()函数,但因对异步编程不熟悉,未能解决问题。
环境信息
- Airflow: 2.5.1
- Python: 3.7
原代码
import discord from tabulate import tabulate import asyncio import pandas as pd async def post_to_discord(df, channel_id, bot_token, as_message=True, num_rows=5): intents = discord.Intents.default() intents.members = True client = discord.Client(intents=intents) try: @client.event async def on_ready(): channel = client.get_channel(channel_id) if as_message: # Post the dataframe as a message, num_rows rows at a time for i in range(0, len(df), num_rows): message = tabulate(df.iloc[i:i+num_rows,:], headers='keys', tablefmt='pipe', showindex=False) await channel.send(message) else: # Send the dataframe as a CSV file df.to_csv("dataframe.csv", index=False) with open("dataframe.csv", "rb") as f: await channel.send(file=discord.File(f)) # client.run(bot_token) await client.start(bot_token) await client.wait_until_ready() finally: await client.close() async def main(df, channel_id, bot_token, as_message=True, num_rows=5): # loop = asyncio.get_event_loop() # result = loop.run_until_complete(post_to_discord(df, channel_id, bot_token, as_message, num_rows)) result = asyncio.run(post_to_discord(df, channel_id, bot_token, as_message, num_rows)) await result return result if __name__ =='__main__': main()
问题分析与解决方案
核心问题点
main函数错误使用asyncio.run():asyncio.run()会直接运行异步函数并返回结果,无需额外await,且在Airflow上下文里手动创建新循环易引发资源泄漏。on_ready事件未主动终止客户端:Discord.py的Client在on_ready执行完成后会持续监听事件,不会自动退出,导致任务卡住无法进入下一个步骤。- Airflow中异步任务运行方式不正确:Airflow默认同步执行任务,需适配异步函数的运行逻辑。
修正后的代码
方式一:调整Discord客户端逻辑,确保任务完成后主动退出
import discord from tabulate import tabulate import asyncio import pandas as pd import os async def post_to_discord(df, channel_id, bot_token, as_message=True, num_rows=5): intents = discord.Intents.default() # 不需要成员权限可关闭,降低权限需求 # intents.members = True client = discord.Client(intents=intents) async def send_data(): channel = client.get_channel(channel_id) if not channel: raise ValueError(f"无法找到ID为{channel_id}的频道") if as_message: # 分批次发送格式化后的DataFrame消息 for i in range(0, len(df), num_rows): chunk = df.iloc[i:i+num_rows, :] message = tabulate(chunk, headers='keys', tablefmt='pipe', showindex=False) await channel.send(message) else: # 发送CSV文件并清理临时文件 csv_path = "temp_dataframe.csv" df.to_csv(csv_path, index=False) try: with open(csv_path, "rb") as f: await channel.send(file=discord.File(f, filename="dataframe.csv")) finally: if os.path.exists(csv_path): os.remove(csv_path) @client.event async def on_ready(): try: await send_data() finally: # 任务完成后主动关闭客户端连接 await client.close() try: await client.start(bot_token) except Exception as e: print(f"发送失败: {str(e)}") raise # 适配Airflow同步任务的运行入口 def run_post_to_discord(df, channel_id, bot_token, as_message=True, num_rows=5): loop = asyncio.get_event_loop() loop.run_until_complete(post_to_discord(df, channel_id, bot_token, as_message, num_rows)) # 本地测试用例 if __name__ == '__main__': test_df = pd.DataFrame({ "Name": ["Alice", "Bob", "Charlie"], "Age": [25, 30, 35] }) run_post_to_discord(test_df, 123456789, "YOUR_BOT_TOKEN", as_message=True)
方式二:使用Airflow的AsyncOperator(推荐)
Airflow 2.x支持直接用AsyncPythonOperator运行异步函数,无需手动处理事件循环:
from airflow import DAG from airflow.operators.python import AsyncPythonOperator from datetime import datetime import pandas as pd # 导入上面修正后的post_to_discord函数 default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), } with DAG('discord_dataframe_dag', default_args=default_args, schedule_interval='@daily') as dag: def get_target_df(): # 替换为你的DataFrame获取逻辑 return pd.DataFrame({ "Name": ["Alice", "Bob", "Charlie"], "Age": [25, 30, 35] }) post_task = AsyncPythonOperator( task_id='post_to_discord', python_callable=post_to_discord, op_kwargs={ 'df': get_target_df(), 'channel_id': 123456789, 'bot_token': "YOUR_BOT_TOKEN", 'as_message': True, 'num_rows': 5 } ) post_task
关键优化说明
- 主动关闭连接:在
on_ready事件完成数据发送后立即调用await client.close(),确保连接正常关闭,任务能顺利结束。 - 临时文件清理:发送CSV文件后删除临时文件,避免Airflow Worker节点出现文件堆积。
- 适配Airflow运行逻辑:同步任务用
asyncio.get_event_loop().run_until_complete()运行异步函数;Airflow 2.x环境下用AsyncPythonOperator更简洁高效。 - 权限精简:无需成员列表权限时关闭
intents.members,减少机器人权限需求,避免不必要的权限报错。
内容的提问来源于stack exchange,提问作者NinjaWarrior
相关产品推荐
相关产品推荐

