Django应用中Python实现的Singleton Pulsar Producer异常排查
在Django应用中实现Apache Pulsar单例生产者,需求是每秒发送数百条消息,复用首次创建的生产者而非每次新建。
问题表现:
- 本地环境运行正常,但部署到AWS ECS生产环境后失效
- 日志显示生产者创建成功,但发送消息时无报错却卡住
- SSH进入Docker容器,通过Python shell导入
PULSAR_PRODUCER时,其值为None
单例生产者实现代码
import logging import os import threading import pulsar from sastaticketpk.settings.config import PULSAR_ENABLED logger = logging.getLogger(__name__) PULSAR_ENV = os.environ.get("PULSAR_CONTAINER") class PulsarClient: __producer = None __lock = threading.Lock() # 确保单例生产者初始化线程安全 def __init__(self): """ 私有构造函数 """ raise RuntimeError("Call get_producer() instead") @classmethod def initialize_producer(cls): if PULSAR_ENABLED and PULSAR_ENV: logger.info("Creating pulsar producer") client = pulsar.Client( "pulsar://k8s-tooling-pulsarpr-7.elb.ap-southeast-1.amazonaws.com:6650" ) # 创建生产者 try: cls.__producer = client.create_producer( "persistent://public/default/gaf", send_timeout_millis=1000 ) logger.info("Producer created successfully") except pulsar._pulsar.TopicNotFound as e: logger.error(f"Error creating producer: {e}") except Exception as e: logger.error(f"Error creating producer: {e}") @classmethod def get_producer(cls): logger.info(f"Producer value {cls.__producer}") if cls.__producer is None: with cls.__lock: if cls.__producer is None: cls.initialize_producer() logger.info(f"Is producer connected {cls.__producer}") return cls.__producer
应用启动时的调用代码
PULSAR_PRODUCER = PulsarClient.get_producer()
环境变量未正确注入:检查ECS任务定义中是否正确配置了
PULSAR_ENABLED和PULSAR_ENV环境变量,进入容器执行echo $PULSAR_ENABLED和echo $PULSAR_CONTAINER确认值是否存在且符合预期。如果任一变量为空,initialize_producer内的创建逻辑不会执行,导致__producer始终为None。网络与安全组限制:AWS ECS环境中,容器所在的安全组、网络ACL可能未开放Pulsar集群ELB的6650端口。虽然日志显示“创建成功”,但可能是客户端初始化未真正完成,或创建生产者时的网络异常被捕获后未处理
__producer为None的情况。建议在initialize_producer中添加客户端连接状态日志,或临时在容器内执行telnet k8s-tooling-pulsarpr-7.elb.ap-southeast-1.amazonaws.com 6650测试连通性。多进程服务器的变量隔离:如果Django使用Gunicorn等多worker模式,启动时在主进程创建的
PULSAR_PRODUCER不会被子进程继承,子进程中该变量会被重置为None。此时应取消启动时的提前初始化,改为在每个进程首次发送消息时调用PulsarClient.get_producer()懒加载(注意:Pulsar生产者非进程安全,多进程场景下每个进程需独立创建生产者实例)。异常处理后的状态遗漏:
initialize_producer中如果创建生产者时抛出异常,cls.__producer仍为None,后续调用get_producer会重复尝试创建。若因网络问题导致每次创建都失败,会持续返回None。建议在异常分支中添加重试逻辑,或标记初始化失败状态避免重复无意义尝试。连接未配置自动重连:Pulsar客户端可能因网络波动断开连接,但当前代码无重连逻辑。可在创建
pulsar.Client时添加connection_timeout、operation_timeout等参数,或在发送消息前检查生产者状态,若已断开则重新初始化。
内容的提问来源于stack exchange,提问作者Usman Hussain

