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

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()

问题分析与解决方案

核心问题点

  1. main函数错误使用asyncio.run():asyncio.run()会直接运行异步函数并返回结果,无需额外await,且在Airflow上下文里手动创建新循环易引发资源泄漏。
  2. on_ready事件未主动终止客户端:Discord.py的Client在on_ready执行完成后会持续监听事件,不会自动退出,导致任务卡住无法进入下一个步骤。
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 13:01:11