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

如何优化异步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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 12:23:19