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

如何用asyncio构建一对多的Twitter API异步处理流水线?

基于asyncio的Twitter API流水线优化方案

核心思路:用异步队列解耦各步骤

把检查、删除、写入三个环节拆成独立的协程,通过asyncio.Queue传递数据,让每个协程只专注自身核心逻辑,调度逻辑由队列统一承接,彻底解决原代码中协程耦合调度的问题。

具体实现代码

import asyncio
from more_itertools import batched

# 保留你已实现的模拟API/工具函数
async def get_tweets(ids):
    # 调用Twitter批量获取推文API
    pass

def check_tweets(responses):
    # 筛选出仍存在的推文ID
    return [resp["id"] for resp in responses if resp.get("exists")]

async def delete_tweet(tweet_id):
    # 调用Twitter删除单条推文API
    pass

def was_deleted(response):
    # 判断删除操作是否成功
    return response.get("deleted") is True

async def write_to_file(tweet_id):
    # 建议用异步文件库aiofiles避免阻塞事件循环
    async with open("deleted_ids.txt", "a") as f:
        await f.write(f"{tweet_id}\n")

# -------------------------- 解耦后的流水线协程 --------------------------
async def batch_check_tweets(input_ids, delete_queue):
    """批量检查推文存在性,将存在的ID送入删除队列"""
    for batch in batched(input_ids, 100):
        try:
            responses = await get_tweets(batch)
            existing_ids = check_tweets(responses)
            for tweet_id in existing_ids:
                await delete_queue.put(tweet_id)
        except Exception as e:
            # 可添加日志、重试等异常处理逻辑
            print(f"批量检查失败: {e}")
    # 所有检查任务完成后,向删除队列发送结束信号
    await delete_queue.put(None)

async def process_deletions(delete_queue, write_queue, max_concurrent=5):
    """从删除队列取ID执行删除,成功的ID送入写入队列"""
    # 用信号量控制并发删除数,避免触发API速率限制
    semaphore = asyncio.Semaphore(max_concurrent)
    
    while True:
        tweet_id = await delete_queue.get()
        if tweet_id is None:
            # 传递结束信号给写入队列
            await write_queue.put(None)
            break
        
        try:
            async with semaphore:
                response = await delete_tweet(tweet_id)
                if was_deleted(response):
                    await write_queue.put(tweet_id)
        except Exception as e:
            # 记录删除失败的ID或添加重试逻辑
            print(f"删除推文{tweet_id}失败: {e}")
        finally:
            delete_queue.task_done()

async def process_writes(write_queue):
    """从写入队列取ID,执行文件写入"""
    while True:
        tweet_id = await write_queue.get()
        if tweet_id is None:
            break
        
        try:
            await write_to_file(tweet_id)
        except Exception as e:
            print(f"写入推文{tweet_id}失败: {e}")
        finally:
            write_queue.task_done()

async def main():
    tweet_ids = [...]  # 你的数千条推文ID列表
    # 创建队列控制数据流转,设置最大容量避免内存占用过高
    delete_queue = asyncio.Queue(maxsize=100)
    write_queue = asyncio.Queue(maxsize=100)
    
    # 启动三个独立的流水线协程
    tasks = [
        asyncio.create_task(batch_check_tweets(tweet_ids, delete_queue)),
        asyncio.create_task(process_deletions(delete_queue, write_queue)),
        asyncio.create_task(process_writes(write_queue))
    ]
    
    # 等待所有任务完成,确保队列中所有数据都被处理
    await asyncio.gather(*tasks)
    await delete_queue.join()
    await write_queue.join()

if __name__ == "__main__":
    asyncio.run(main())

方案优势

  • 职责单一:每个协程只负责一个环节,比如batch_check_tweets仅处理批量检查和入队,逻辑清晰,后续修改不会影响其他模块。
  • 易于扩展:添加异常重试、日志监控、任务取消等功能时,只需在对应协程内部调整,无需改动整个流水线逻辑。
  • 速率可控:通过asyncio.Semaphore可以精准控制删除操作的并发数,避免触发Twitter API的速率限制。
  • 任务可追踪:利用队列的task_done()和join()方法,能准确跟踪所有任务的完成状态,实现优雅退出。

内容的提问来源于stack exchange,提问作者Michael Banks

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 19:58:19