如何用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
相关产品推荐
相关产品推荐

