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

如何让Python Pulsar消费者在Docker环境中等待依赖服务启动

解决方案

直接在Python脚本中增加服务可用性检测重试逻辑,捕获Pulsar连接异常,重试直到依赖服务正常启动后再执行后续消费逻辑,修改后的代码如下:

import pulsar
import time

def wait_for_pulsar_service(pulsar_url, max_retries=30, retry_interval=2):
    """等待Pulsar服务可用"""
    retries = 0
    while retries < max_retries:
        try:
            # 尝试创建测试客户端检测连接可用性
            test_client = pulsar.Client(pulsar_url, operation_timeout_seconds=3)
            test_client.close()
            print("Pulsar服务已正常启动,开始初始化消费者")
            return True
        except Exception as e:
            retries += 1
            print(f"Pulsar服务未就绪,第{retries}次重试,错误信息:{str(e)}")
            time.sleep(retry_interval)
    print(f"重试{max_retries}次后Pulsar服务仍不可用,退出程序")
    return False

def initialize_consumer():
    pulsar_url = 'pulsar://localhost:6650'
    # 先等待Pulsar服务就绪再执行后续逻辑
    if not wait_for_pulsar_service(pulsar_url):
        exit(1)
    
    client = pulsar.Client(pulsar_url)
    consumer = client.subscribe('my-topic', 'my-subscription')

    while True:
        msg = consumer.receive()
        try:
            output_string = f"收到消息 {msg.data()} 消息ID={msg.message_id()}"
            print(output_string)
            with open('./output.txt', 'a') as f:
                f.write(output_string + '\n')
            # 确认消息处理成功
            consumer.acknowledge(msg)
        except:
            # 消息处理失败标记
            consumer.negative_acknowledge(msg)

    client.close()

if __name__ == "__main__":
    initialize_consumer()
  • 可根据实际服务启动时长调整max_retries(最大重试次数)和retry_interval(每次重试间隔秒数)参数
  • 如果使用docker-compose编排容器,可配合服务健康检查规则进一步保证启动顺序,减少不必要的重试等待

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 03:06:03