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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 14:08:14