【发布时间】:2023-02-12 02:19:11
【问题描述】:
我正在尝试将数据框发布到 Discord 频道。但是,我在让 Discord.py 关闭连接并继续执行下一个任务时遇到问题。我已经尝试使用此线程 (How to run async function in Airflow?) 中建议的事件循环以及 asyncio.run() 函数。不太熟悉异步,希望能在这里得到一些指示。下面是我尝试在 DAG 和 Task 中导入但没有成功的 Python 代码。提前致谢!
气流:2.5.1
蟒蛇: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()
【问题讨论】:
-
你为什么使用异步?让Airflow控制任务的并行执行不是更容易吗?
-
当你使用
loop.run_until_complete()/asyncio.run()时删除await result。此外,将async def main更改为def main。 -
@SultanOrazbayev 如果我不使用异步,则任务将成功而无需将消息发布到 Discord。它不会等待连接建立。
-
@aaron 感谢您的建议。进行了这两项更改(def main 和删除等待结果),但在 Discord 中发布消息后任务继续运行(不关闭连接)。
-
它卡在
await client.wait_until_ready()了吗?
标签: python pandas async-await airflow