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
相关产品推荐
相关产品推荐

