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

在Django中自动运行Kafka Consumer的实现问题

Django Kafka消费者后台运行方案

方法1:自定义Django管理命令 + 进程管理工具(推荐生产环境)

步骤1:编写自定义消费者命令

在消费者Django项目的app/management/commands目录下创建start_kafka_consumer.py:

from django.core.management.base import BaseCommand
from kafka import KafkaConsumer
import json

class Command(BaseCommand):
    help = '启动Kafka消费者'

    def handle(self, *args, **options):
        consumer = KafkaConsumer(
            'your_topic_name',
            bootstrap_servers=['your_kafka_host:9092'],
            auto_offset_reset='earliest',
            value_deserializer=lambda m: json.loads(m.decode('utf-8'))
        )
        for message in consumer:
            # 替换为你的消息处理逻辑
            print(f"收到消息: {message.value}")

步骤2:用Supervisor管理两个进程

安装Supervisor后,创建配置文件/etc/supervisor/conf.d/django_kafka.conf:

[program:django_server]
command=/path/to/venv/bin/gunicorn your_project.wsgi:application --bind 0.0.0.0:8000
directory=/path/to/your/django/project
user=ubuntu
autostart=true
autorestart=true
stdout_logfile=/var/log/django_server.log
stderr_logfile=/var/log/django_server_err.log

[program:kafka_consumer]
command=/path/to/venv/bin/python manage.py start_kafka_consumer
directory=/path/to/your/django/project
user=ubuntu
autostart=true
autorestart=true
stdout_logfile=/var/log/kafka_consumer.log
stderr_logfile=/var/log/kafka_consumer_err.log

执行命令更新配置并启动:

supervisorctl reread
supervisorctl update
supervisorctl start all

方法2:在Django启动时启动后台线程(适合测试/轻量场景)

修改消费者项目的wsgi.py(或asgi.py),添加后台线程启动消费者:

import threading
from kafka import KafkaConsumer
import json

def run_consumer():
    consumer = KafkaConsumer(
        'your_topic_name',
        bootstrap_servers=['your_kafka_host:9092'],
        auto_offset_reset='earliest',
        value_deserializer=lambda m: json.loads(m.decode('utf-8'))
    )
    for message in consumer:
        # 替换为你的消息处理逻辑
        print(f"收到消息: {message.value}")

# 启动消费者守护线程
threading.Thread(target=run_consumer, daemon=True).start()

# 原有WSGI配置
import os
from django.core.wsgi import get_wsgi_application

os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'your_project.settings')
application = get_wsgi_application()

注意:生产环境用Gunicorn开启多worker时,会启动多个消费者实例,需根据需求调整worker数量或用进程锁避免重复消费。

方法3:独立脚本 + Systemd服务

把消费者逻辑写成独立脚本kafka_consumer.py:

from kafka import KafkaConsumer
import json

if __name__ == "__main__":
    consumer = KafkaConsumer(
        'your_topic_name',
        bootstrap_servers=['your_kafka_host:9092'],
        auto_offset_reset='earliest',
        value_deserializer=lambda m: json.loads(m.decode('utf-8'))
    )
    for message in consumer:
        # 替换为你的消息处理逻辑
        print(f"收到消息: {message.value}")

创建Systemd服务文件/etc/systemd/system/kafka-consumer.service:

[Unit]
Description=Kafka Consumer Service
After=network.target

[Service]
User=ubuntu
WorkingDirectory=/path/to/your/project
ExecStart=/path/to/venv/bin/python kafka_consumer.py
Restart=always

[Install]
WantedBy=multi-user.target

启动并设置开机自启:

systemctl daemon-reload
systemctl start kafka-consumer.service
systemctl enable kafka-consumer.service

同时用Systemd单独管理Django服务,确保两者独立运行、互不阻塞。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 19:12:19