Python pulsar-client 的 python 异步支持

来源:建站教程作者:台湾程序员头衔:程序员
导读:本期聚焦于台湾程序员创作的《Python pulsar-client 的 python 异步支持》,敬请观看详情。Apache Pulsar 作为新一代分布式消息队列,官方提供了 Python 客户端 pulsar-client。很多使用者在 FastAPI、aiohttp 这类基于 asyncio 的框架里接入 Pulsar 时,会直接使用同步的 Client 和 Producer,结果发现事件循环被阻塞,接口响应变慢。实际上,从 pulsar-client 2.10 版本开始,官方已经原生支持 Python 异步接口,也就是 pulsar_

Apache Pulsar 作为新一代分布式消息队列,官方提供了 Python 客户端 pulsar-client。很多使用者在 FastAPI、aiohttp 这类基于 asyncio 的框架里接入 Pulsar 时,会直接使用同步的 ClientProducer,结果发现事件循环被阻塞,接口响应变慢。实际上,从 pulsar-client 2.10 版本开始,官方已经原生支持 Python 异步接口,也就是 pulsar_async 模块,可以直接在协程中完成消息的生产和消费,无需借助线程池包装。本文将系统介绍这套异步接口的原理和用法。

Python pulsar-client 的 python 异步支持

为什么需要异步接口:同步客户端在协程中的坑

先看同步接口的问题。pulsar-client 的同步 producer.send() 在内部会等待服务端的确认(取决于发送模式),consumer.receive() 则会一直阻塞直到有消息到达。如果把这些调用直接写在协程里,整个事件循环都会被卡住,同一进程内的其他协程、HTTP 请求全部无法调度,表现出来就是接口偶发性地整体卡死。

常见的绕过方案是用 asyncio.to_thread 或线程池把同步调用丢到子线程执行。这种方式能解决阻塞问题,但代价不小:每个阻塞调用占用一个线程,消息量大时线程切换开销明显;消费逻辑如果需要持续循环,还要自己维护线程生命周期;此外,跨线程传递返回值和异常会让代码变得繁琐。相比之下,原生异步接口直接返回 awaitable 对象,调度完全交给事件循环,资源占用和代码可读性都更好。

需要注意版本问题。异步支持从 pulsar-client 2.10.0 开始引入,早期版本没有 AsyncProducerAsyncConsumer。安装前先确认版本:

pip install pulsar-client==2.10.1
# 如果需要认证和 avro 支持
pip install pulsar-client[all]==2.10.1

异步接口的核心用法:AsyncProducer 与 AsyncConsumer

异步接口的入口同样是 Client,但生产者和消费者要显式指定 async_ 参数。由于 Python 中 async 是保留字,官方 API 使用了带下划线的形式,这是初学者最容易忽略的一个细节。

创建异步生产者并发送消息的完整示例如下:

import asyncio
import pulsar

async def main():
    client = pulsar.Client('pulsar://127.0.0.1:6650')

    # 关键参数:async_=True 返回异步生产者
    producer = await client.create_producer(
        'persistent://public/default/my-topic',
        async_=True,
        batching_enabled=True,
        batching_max_publish_delay_ms=10
    )

    # send 返回 asyncio.Future,可以 await 拿到消息 id
    msg_id = await producer.send(b'hello async pulsar')
    print('消息已确认, id:', msg_id)

    await producer.close()
    client.close()

asyncio.run(main())

这段代码中有几点值得展开。第一,create_producer 本身也变成了异步操作,需要 await。第二,send() 返回的是 asyncio.Future,await 它可以拿到服务端返回的 MessageId,确保消息已持久化;如果只是尽快发出去、不关心单条确认,也可以不 await 直接依赖批量机制。第三,开启了 batching_enabled 后,客户端会在本地攒一批消息再发送,通常能显著提升吞吐。

异步消费者的用法类似,重点在于它提供了回调式的消息处理机制,通过 listener 参数注册一个异步函数,每收到一条消息就触发一次:

import asyncio
import pulsar

async def handle_message(consumer, message):
    # 处理业务逻辑,这里可以是数据库操作、HTTP 调用等异步任务
    print('收到消息:', message.data().decode('utf-8'))
    await asyncio.sleep(0.01)
    # 处理成功后再确认
    consumer.acknowledge(message)

async def main():
    client = pulsar.Client('pulsar://127.0.0.1:6650')
    consumer = await client.subscribe(
        'persistent://public/default/my-topic',
        subscription_name='my-sub',
        async_=True,
        consumer_type=pulsar.ConsumerType.Shared,
        message_listener=handle_message
    )
    # 保持主协程存活
    await asyncio.Event().wait()
    await consumer.close()
    client.close()

asyncio.run(main())

监听器模式下不需要自己写 receive 循环,客户端内部会在事件循环上调度回调,处理完手动调用 acknowledge 即可。如果处理失败,调用 negative_acknowledge 让消息稍后重投,比直接抛异常更可控。

同步与异步模式的对比及选型建议

两种模式各有适用场景,简单对比一下:

维度同步接口异步接口
运行环境普通脚本、多线程程序asyncio 事件循环、FastAPI 等异步框架
阻塞行为send 和 receive 会阻塞当前线程全程不阻塞事件循环
并发能力依赖多线程或多进程单线程内天然高并发
消费模式while 循环调用 receivelistener 回调或异步 receive
版本要求所有版本2.10.0 及以上

如果项目主体是普通脚本或者 Django 这类同步框架,直接用同步客户端更简单直接;如果是 FastAPI、Sanic、aiohttp 项目,或者需要在一个进程里同时维护成百上千个主题的收发,异步接口是更合理的选择。还有一种混合场景:遗留代码里已经有同步消费者,新模块用异步,可以在同一个进程中共存两个 Client 实例,但要注意分别关闭。

常见配置项与踩坑提示

连接参数方面,生产环境建议显式设置 operation_timeout_seconds(单次操作超时)和 connection_timeout_ms,避免网络抖动时调用长时间挂起。认证集群使用 authentication 参数传入 Token 或 TLS 证书配置,异步接口的用法与同步版完全一致。

client = pulsar.Client(
    'pulsar+ssl://broker.ippipp.com:6651',
    authentication=pulsar.AuthenticationToken('your-token'),
    operation_timeout_seconds=30,
    connection_timeout_ms=10000
)

几个容易踩的坑需要提醒。一是忘记关闭资源:producer、consumer 和 client 都要调用 close(),建议用 try/finally 或异步上下文管理包裹,否则连接泄漏在长时间运行的服务里会逐渐耗尽句柄。二是 async_ 参数漏写,此时返回的是同步对象,调用 send 得到的不是 Future,await 它会直接报错。三是监听器里不要执行 CPU 密集型任务,事件循环被计算占满同样会导致其他协程卡顿,重计算应该交给 loop.run_in_executor 或独立进程。

最后是优雅退出问题。服务收到 SIGTERM 时,应该在关闭 consumer 之前停止接收新消息、等待在途消息处理完并调用 unsubscribeclose,配合 Shared 订阅模式可以让其他消费者实例快速接管分区,实现平滑发布。

修改时间:2026-09-12 00:32:39

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/20260912/54997.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。