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

如何通过Celery REST API实现带表单参数的HTTP POST任务队列及重试?

Absolutely! Celery is tailor-made for this kind of asynchronous task handling with retry logic. Let’s walk through exactly how to build this solution step by step.


1. Define the Celery Task with Retry Logic

First, set up your Celery instance and create a task that handles the POST request, checks for non-200 responses, and retries on failure. We’ll use requests for HTTP calls and configure retry parameters directly in the task.

Step 1.1: Install Dependencies

pip install celery requests redis  # Using Redis as broker/backend; swap for RabbitMQ if preferred

Step 1.2: Write the Celery Task

Create a file (e.g., tasks.py):

from celery import Celery
import requests

# Initialize Celery with your broker/backend
app = Celery(
    "post_tasks",
    broker="redis://localhost:6379/0",
    backend="redis://localhost:6379/0"
)

@app.task(bind=True, max_retries=3)  # Max 3 retries total
def send_post_request(self, target_url, form_data):
    try:
        # Send POST request with form data (application/x-www-form-urlencoded)
        response = requests.post(target_url, data=form_data)
        # Raise an exception if status code is non-200
        response.raise_for_status()
        
        return {
            "status": "success",
            "response_code": response.status_code,
            "content": response.text[:500]  # Truncate long content for brevity
        }
    except requests.exceptions.HTTPError as exc:
        # Retry on non-200 responses (400, 404, etc.)
        print(f"Request failed with {exc}, retrying...")
        self.retry(exc=exc, countdown=60)  # Wait 60 seconds before retrying
    except requests.exceptions.RequestException as exc:
        # Handle network errors (timeouts, DNS failures) with shorter retry interval
        print(f"Network error: {exc}, retrying sooner...")
        self.retry(exc=exc, countdown=15)

Key details here:

  • bind=True lets the task access its own instance (self) to call self.retry()
  • max_retries sets the total number of retry attempts
  • countdown controls how long to wait before retrying (adjust based on your needs)

2. Build a REST API to Submit Tasks

Next, create a simple REST endpoint (we’ll use Flask here) to accept task parameters and queue them via Celery.

Step 2.1: Write the Flask API

Create api.py:

from flask import Flask, request, jsonify
from tasks import send_post_request

app = Flask(__name__)

@app.route("/submit-post-task", methods=["POST"])
def submit_task():
    # Parse incoming JSON payload
    payload = request.get_json()
    
    # Validate required fields
    required_fields = ["target_url", "form_data"]
    if not all(field in payload for field in required_fields):
        return jsonify({
            "error": f"Missing required fields: {', '.join(required_fields)}"
        }), 400
    
    # Queue the Celery task
    task = send_post_request.delay(
        target_url=payload["target_url"],
        form_data=payload["form_data"]
    )
    
    # Return task ID for status checking
    return jsonify({
        "task_id": task.id,
        "status": "task queued"
    }), 202

@app.route("/task-status/<task_id>", methods=["GET"])
def get_task_status(task_id):
    task = send_post_request.AsyncResult(task_id)
    
    if task.state == "PENDING":
        response = {"state": "pending", "message": "Task is waiting to run"}
    elif task.state == "SUCCESS":
        response = {"state": "success", "result": task.result}
    elif task.state == "RETRY":
        response = {"state": "retrying", "message": f"Next attempt in {task.next_retry}s"}
    else:
        # Handle failures (after max retries)
        response = {"state": "failed", "error": str(task.info)}
    
    return jsonify(response)

if __name__ == "__main__":
    app.run(host="0.0.0.0", port=5000)

3. Test the Workflow

Step 3.1: Start Celery Worker

In one terminal:

celery -A tasks worker --loglevel=info

Step 3.2: Start the Flask API

In another terminal:

python api.py

Step 3.3: Submit a Task via API

Use curl or Postman to send a request:

curl -X POST -H "Content-Type: application/json" -d '{
    "target_url": "https://external-api.example.com/your-endpoint",
    "form_data": {"username": "test-user", "action": "submit-form"}
}' http://localhost:5000/submit-post-task

You’ll get a response like:

{"task_id": "abc123def456", "status": "task queued"}

Step 3.4: Check Task Status

curl http://localhost:5000/task-status/abc123def456

4. Customization Tips
  • Adjust max_retries and countdown to match your external API’s rate limits or recovery time
  • Use retry_backoff instead of countdown for exponential backoff (e.g., self.retry(exc=exc, retry_backoff=True))
  • Add logging instead of print statements for production use
  • For larger payloads, consider passing references to stored data instead of sending raw form data in the API request

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:32:36