如何实现Telegram机器人与Twitter交互程序的aioprocessing数据传输?
需求说明
我有一个Telegram机器人和一个调用Twitter API的独立程序,二者部署在同一服务器的同一目录下。需求是:当用户向机器人处理器提交数据时,需通过aioprocessing将数据传输至另一程序并加入其队列。
MYPROGRAM.PY
import logging import asyncio from datetime import datetime, timedelta import httpx import aioprocessing logging.basicConfig(level=logging.INFO) task_pool = [] account_queue = aioprocessing.AioQueue() class Account_Data: def __init__(self, auth_token, proxy_ip, proxy_port, proxy_login, proxy_pass, bearer_token, x_csrf_token, user_agent, nickname, conversations, gif_number, text_message, messages_in_second, work_durations, break_cycle): self.auth_token = auth_token self.proxy_ip = proxy_ip self.proxy_port = proxy_port self.proxy_login = proxy_login self.proxy_pass = proxy_pass self.bearer_token = bearer_token self.x_csrf_token = x_csrf_token self.user_agent = user_agent self.nickname = nickname self.conversations = conversations self.gif_number = gif_number self.text_message = text_message self.messages_in_second = messages_in_second self.work_durations = work_durations self.break_cycle = break_cycle async def add_account_task(account): # 检查该账号是否已有对应的任务 if not any(task.get_name() == account.nickname for task in task_pool): logging.info(f"为账号添加任务: {account.nickname}") task = asyncio.create_task(process_chat(account), name=account.nickname) task_pool.append(task) else: logging.info(f"账号 {account.nickname} 的任务已存在") async def process_queue(): while True: # 检查队列中是否有新账号需要添加 if not account_queue.empty(): logging.info("从队列获取新数据") account_data = await account_queue.coro_get() # 从队列获取账号数据 await add_account_task(account_data) # 为账号添加任务 await asyncio.sleep(1) # 短暂延迟,避免持续轮询 async def work_with_chats(accounts): accounts_data_list = [] for nickname in accounts: account_data = await get_user_retweet_data(nickname) if account_data: account_obj = Account_Data( auth_token=account_data['auth_token'], proxy_ip=account_data['proxy_ip'], proxy_port=account_data['proxy_port'], proxy_login=account_data['proxy_login'], proxy_pass=account_data['proxy_pass'], bearer_token=account_data['bearer_token'], x_csrf_token=account_data['x_csrf_token'], user_agent=account_data['user_agent'], nickname=nickname, conversations = list(filter(lambda x: len(x.split("-")) == 1, account_data['conversations'])), gif_number=account_data['gif_number'], text_message=account_data['text_message'], messages_in_second=account_data['messages_in_second'], work_durations=account_data["work_duration"], break_cycle = False ) accounts_data_list.append(account_obj) else: logging.error(f"无法提取账号 {nickname} 的数据") for account in accounts_data_list: await add_account_task(account) # 为每个账号添加任务 async def work_with_gif(account, conversation_id, session): ................... async def ReTwit(account, conversation_id, session): ........................................... async def sending_message(account, conversation_id, session): .............................................. async def process_chat(account): end_time = datetime.now() + timedelta(seconds=account.work_durations) while True: if account.break_cycle is True: logging.info(f"删除账号 {account.nickname} 的任务") task_pool[:] = [task for task in task_pool if task.get_name() != account.nickname] break for conversation_id in account.conversations: if datetime.now() >= end_time: # await message.answer(f"账号 {account.nickname} 的群发流程已超时结束") account.break_cycle = True # 此处应为赋值而非比较 break proxies = f"socks5://{account.proxy_login}:{account.proxy_pass}@{account.proxy_ip}:{account.proxy_port}" async with httpx.AsyncClient(proxies=proxies) as session: await ReTwit(account, conversation_id, session) await work_with_gif(account, conversation_id, session) await sending_message(account, conversation_id, session) await asyncio.sleep(3600 / account.messages_in_second) async def task_manager(): while True: if task_pool: logging.info(f"当前执行 {len(task_pool)} 个任务。") else: logging.info("任务池为空,等待新任务。") await asyncio.sleep(5) async def stop_account_task(nickname): """根据昵称停止群发任务""" for task in task_pool: if task.get_name() == nickname: logging.info(f"停止账号 {nickname} 的任务") task.cancel() # 取消任务 try: if not task.done(): await task # 等待任务完成(如果尚未结束) except asyncio.CancelledError: logging.info(f"账号 {nickname} 的任务已被取消") # 取消后从任务池中移除 task_pool.remove(task) logging.info(f"账号 {nickname} 的任务已从任务池删除") break else: logging.warning(f"未找到账号 {nickname} 的任务") async def main(): await asyncio.gather(task_manager(), process_queue()) # 同时执行任务管理和队列处理 if __name__ == "__main__": asyncio.run(main())
MESSAGE_HANDLER.PY
async def start_program(self, callback_query: CallbackQuery, state: FSMContext): user_id = callback_query.from_user.id # 检查当前用户是否已选择昵称 if user_id in self.user_selected_nicknames and self.user_selected_nicknames[user_id]: selected_nicknames = self.user_selected_nicknames[user_id] await callback_query.message.answer(f"已启动以下昵称的群发流程: {', '.join(selected_nicknames)}") # 调用work_with_chats处理当前用户的昵称,加入队列 await self.task_queue.coro_put(list(selected_nicknames)) # 完成后清空昵称列表 self.user_selected_nicknames[user_id].clear() await state.clear() else: await callback_query.message.answer("您未选择任何昵称。")
内容的提问来源于stack exchange,提问作者user25074879
相关产品推荐
相关产品推荐

