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

如何在Celery任务中监听Django Web API的请求?

需求可行性分析与实现方案

你的需求完全可行,但不能直接在Celery任务内部监听Post API请求——因为Celery任务是周期性/异步执行的独立进程,和Django的Web服务进程运行在不同上下文,无法直接捕获HTTP请求。换个思路:让API被请求时留下记录,再让Celery任务定期检查这些记录即可。

具体实现方案

1. 记录API请求轨迹

在你的Post API视图中,每次收到请求时,把请求信息存入数据库或缓存(推荐Redis,性能更高):

方案1:用数据库记录(适合低并发场景)

先创建日志模型:

# models.py
from django.db import models
from django.utils import timezone

class APIRequestLog(models.Model):
    path = models.CharField(max_length=255, verbose_name="API路径")
    method = models.CharField(max_length=10, verbose_name="请求方法")
    timestamp = models.DateTimeField(default=timezone.now, verbose_name="请求时间")
    status_code = models.IntegerField(verbose_name="响应状态码")

    class Meta:
        verbose_name = "API请求日志"
        verbose_name_plural = "API请求日志"

然后在Post API视图中写入日志:

# views.py
from django.http import JsonResponse
from .models import APIRequestLog

def post_api_view(request):
    # 你的API业务逻辑
    response = JsonResponse({"status": "ok"})
    
    # 写入请求日志
    APIRequestLog.objects.create(
        path=request.path,
        method=request.method,
        status_code=response.status_code
    )
    return response

方案2:用Redis记录(适合高并发场景)

Redis的有序集合或计数器能高效处理高频请求记录:

# views.py
import redis
from django.http import JsonResponse
from django.utils import timezone

# 初始化Redis连接(建议用Django-redis库统一管理)
r = redis.Redis(host="localhost", port=6379, db=0)

def post_api_view(request):
    # 你的API业务逻辑
    response = JsonResponse({"status": "ok"})
    
    # 把请求时间戳存入有序集合
    r.zadd("post_api_requests", {timezone.now().timestamp(): 1})
    return response

2. 修改Celery任务检查记录

让Celery任务在每次执行时,查询指定时间窗口内的请求记录,判断是否有Post API请求:

对应数据库记录的任务实现

# tasks.py
from celery import shared_task
from .models import APIRequestLog
from django.utils import timezone
from datetime import timedelta

@shared_task(bind=True)
def my_task(self):
    # 根据Celery Beat的调度周期设置检查窗口(比如每5分钟执行一次,就查最近5分钟的记录)
    check_window = timezone.now() - timedelta(minutes=5)
    
    # 查询是否有目标Post API的请求
    has_request = APIRequestLog.objects.filter(
        method="POST",
        path="/your-post-api-path/",  # 替换成你的实际API路径
        timestamp__gte=check_window
    ).exists()

    if has_request:
        # 检测到请求后的处理逻辑
        return "检测到Post API请求"
    else:
        # 未检测到请求的处理逻辑
        return "未检测到Post API请求"

对应Redis记录的任务实现

# tasks.py
from celery import shared_task
import redis
from django.utils import timezone
from datetime import timedelta

r = redis.Redis(host="localhost", port=6379, db=0)

@shared_task(bind=True)
def my_task(self):
    check_window = timezone.now() - timedelta(minutes=5)
    # 统计时间窗口内的请求数量
    request_count = r.zcount("post_api_requests", check_window.timestamp(), "+inf")
    
    if request_count > 0:
        # 清理过期记录,避免Redis占用过多空间
        r.zremrangebyscore("post_api_requests", "-inf", check_window.timestamp())
        return f"检测到{request_count}次Post API请求"
    else:
        return "未检测到Post API请求"

关键说明

Celery worker和Django Web服务是独立的进程,无法直接共享HTTP请求上下文,必须通过中间存储(数据库、缓存)传递请求状态。这种间接监听的方式是最合理的实现路径。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 08:30:49