如何通过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.
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=Truelets the task access its own instance (self) to callself.retry()max_retriessets the total number of retry attemptscountdowncontrols how long to wait before retrying (adjust based on your needs)
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)
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
- Adjust
max_retriesandcountdownto match your external API’s rate limits or recovery time - Use
retry_backoffinstead ofcountdownfor 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

