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

如何实现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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 08:04:53