如何优化异步SNS消息发布性能?实现高效发送即遗忘
问题
我尝试以异步方式向AWS SNS批量发送消息,希望实现“发送即遗忘”模式,无需等待publish接口的响应。参考了“sleepy-bois”方案,期望让任务并行执行而非串行。当仅保留send_item()中的sleep逻辑时,性能表现良好:Process time difference为11.07%,Absolute time difference为103.85%;但加入异步SNS publish调用后,这两个指标飙升至1138.72%和1145.18%。请问是否有更优的异步SNS发布实现方式,能让Absolute time difference值接近100%?
当前实现代码如下:
import asyncio import time from aiobotocore.session import get_session sns_session = get_session() async def send_item(destination, message, delay_ms): delay_seconds = delay_ms / 1000 await asyncio.sleep(delay_seconds) async with sns_session.create_client('sns') as sns_client: await sns_client.publish( TopicArn=destination, Message=message, ) async def send_items(my_items_list): await asyncio.gather(*[send_item(item["destination"], item["message"], item["waitTime_ms"]) for item in my_items_list]) print("Iterated through items.") return my_items_list[-1]["waitTime"] if __name__ == "__main__": my_list = get_list() start_process_time_ns = time.process_time_ns() time_before_send_ns = time.perf_counter_ns() last_wait_time_ms = asyncio.run(send_items(my_list)) time_after__send_ns = time.perf_counter_ns() end_process_time_ns = time.process_time_ns() total_time_ns = time_after__send_ns - time_before_send_ns process_time_ns = end_process_time_ns - start_process_time_ns expected_time_ns = last_wait_time_ms * 1000000 print(f"Expected time (ns): {expected_time_ns}") print(f"Process time (ns): {process_time_ns}") print(f"Absolute time (ns): {total_time_ns}") print(f"Process time difference: {(process_time_ns/expected_time_ns)*100:.2f}%") print(f"Absolute time difference: {(total_time_ns/expected_time_ns)*100:.2f}%") print("Complete")
优化方案
1. 复用SNS客户端,消除重复连接开销
原代码中每个send_item都会新建SNS客户端,频繁的连接创建与销毁是性能骤降的核心原因。改为在send_items中创建单个客户端,所有任务共享该客户端,大幅减少连接相关的耗时。
2. 实现真正的“发送即遗忘”模式
不需要等待publish接口响应时,可将publish调用封装为后台任务,让事件循环直接调度执行,无需阻塞主流程。若需处理极端发送失败场景,可添加基础异常捕获(但不阻塞主任务)。
3. 控制并发数量
AWS SNS存在请求速率限制,过高并发会触发限流反而降低性能。通过asyncio.Semaphore控制并发任务数,既能避免限流,也能防止本地资源耗尽。
修改后的代码
import asyncio import time from aiobotocore.session import get_session sns_session = get_session() # 根据自身AWS SNS配额调整,例如设为100 MAX_CONCURRENCY = 100 async def send_item(sns_client, destination, message, delay_ms, semaphore): delay_seconds = delay_ms / 1000 await asyncio.sleep(delay_seconds) async with semaphore: # 发送即遗忘:提交后台任务,不等待响应 asyncio.create_task( sns_client.publish( TopicArn=destination, Message=message, ), name=f"sns-publish-{destination}" ) async def send_items(my_items_list): # 创建单个SNS客户端,所有任务共享 async with sns_session.create_client('sns') as sns_client: semaphore = asyncio.Semaphore(MAX_CONCURRENCY) tasks = [ send_item(sns_client, item["destination"], item["message"], item["waitTime_ms"], semaphore) for item in my_items_list ] await asyncio.gather(*tasks) print("Iterated through items.") return my_items_list[-1]["waitTime"] if __name__ == "__main__": my_list = get_list() start_process_time_ns = time.process_time_ns() time_before_send_ns = time.perf_counter_ns() last_wait_time_ms = asyncio.run(send_items(my_list)) time_after__send_ns = time.perf_counter_ns() end_process_time_ns = time.process_time_ns() total_time_ns = time_after__send_ns - time_before_send_ns process_time_ns = end_process_time_ns - start_process_time_ns expected_time_ns = last_wait_time_ms * 1000000 print(f"Expected time (ns): {expected_time_ns}") print(f"Process time (ns): {process_time_ns}") print(f"Absolute time (ns): {total_time_ns}") print(f"Process time difference: {(process_time_ns/expected_time_ns)*100:.2f}%") print(f"Absolute time difference: {(total_time_ns/expected_time_ns)*100:.2f}%") print("Complete")
额外说明
- aiobotocore客户端本身是异步安全的,多个任务共享不会有线程安全问题,且能复用连接池。
- 并发数需参考自身AWS账号的SNS请求配额调整,避免触发限流。
- 若需确保消息至少提交到SNS(非完全“发送即遗忘”),可将
asyncio.create_task改为await,复用客户端后性能仍远优于原代码。
内容的提问来源于stack exchange,提问作者wsa225
相关产品推荐
相关产品推荐

