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

Gunicorn+Supervisor部署下Kafka Producer无法发消息求助

Kafka Producer在Gunicorn+Supervisor部署环境下无法通过业务流程发送消息

问题现象

  • 本地使用runserver运行时,MyModel的post_save信号触发的Kafka消息发送完全正常
  • 在GCP容器中手动调用produce_message函数,消息能正常发送到Kafka Topic并被消费者接收
  • 仅在Gunicorn+Supervisor部署的生产环境下,通过业务流程创建/更新MyModel实例时,消息无法发送到Kafka
  • 已确认全局Producer能正常连接Kafka Broker,且post_save信号已被触发

核心代码实现

全局Producer初始化(settings.py)

def init_producer(broker_url, delivery_timeout, retry_interval):
    try:
        return Producer(
            {
                "bootstrap.servers": broker_url,
                "delivery.timeout.ms": delivery_timeout,
                "request.timeout.ms": retry_interval,
            }
        )
    except Exception as exception:
        logger.error(
            f"Couldn't connect to broker {broker_url} because of exception {exception}"
        )

producer = init_producer(
    KAFKA_BROKER_URL, KAFKA_DELIVERY_TIMEOUT, KAFKA_RETRY_INTERVAL
)

post_save信号与消息发送逻辑

@receiver(post_save, sender=MyModel)
def update_report(sender, instance, **kwargs):
    produce_message(instance.id)

def produce_message(instance_id):
    try:
        data = json.dumps(
            {
                "id": instance_id,
                "client": CLIENT,
                "database": "NAME",
            }
        )
        producer.produce(
            KAFKA_REPORT_TOPIC,
            data.encode("utf-8"),
            on_delivery=report_transaction_handler,
        )
        producer.poll(1)
    except Exception as error:
        logger.log(
            "MONITORING",
            f"Error {error} on producing message for Instance {instance_id}"
        )

投递回调函数

def report_model_handler(err, msg):
    """Called once for each message produced to indicate delivery result.
    Triggered by poll() or flush()."""
    if err is not None:
        value = simplejson.loads(msg.value().decode())
        logger.log(
            "MONITORING",
            f"Couldn't produce message {msg.value()} to topic {KAFKA_REPORT_TOPIC}"
        )
    else:
        logger.info(f"Message delivered to {msg.topic()} [{msg.partition()}]")

部署配置

Supervisor配置

[program:django]
command=gunicorn --config settings/gunicorn.py settings.wsgi
stopasgroup=true
killasgroup=true
autostart=true
autorestart=true
stdout_logfile=/dev/stdout
stdout_logfile_maxbytes=0
stderr_logfile=/dev/stderr
stderr_logfile_maxbytes=0

Gunicorn配置

# Configuration file for gunicorn.
# See http://docs.gunicorn.org/en/latest/configure.html for details.
import multiprocessing
from os import environ

port = environ.get("APP_PORT", 80)
bind = ["0.0.0.0:{}".format(port)]
limit_request_line = 64 * 2**10  # 16KB
workers = environ.get("APP_WORKERS", (2 * multiprocessing.cpu_count()) + 1)
worker_class = "gthread"
timeout = 120
graceful_timeout = 120
preload_app = True
worker_tmp_dir = "/dev/shm"
max_requests = 1000
max_requests_jitter = 50
loglevel = "error"

已确认的排查点

  • 全局Producer实例能成功连接到Kafka Broker
  • post_save信号在业务流程中已被正常触发
  • 手动调用produce_message函数时,消息发送流程完全正常

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 08:54:53