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

同步SOA架构下快慢Worker通信方案及REST RPC集成咨询

Hey there, let's work through your scenario step by step—you've got a mix of synchronous REST SOA, slow/serial workers, firewall constraints, and a Python/Flask/Gunicorn stack, so we'll tailor solutions to fit all those pieces.


First, ditch pure synchronous REST for your slow tasks (30-second latency is way too long for blocking requests). Instead, adopt an asynchronous task processing flow with a message broker, paired with a sync-friendly way for consumers to retrieve results.

For your serial processing requirement (license-limited tasks), RabbitMQ is a better fit than Kafka out of the box:

  • RabbitMQ lets you enforce single-consumer queues easily, guaranteeing tasks are processed one at a time.
  • Kafka can handle serial processing too, but it requires locking tasks to a single partition and limiting consumer group size—more configuration overhead.

Mixing Sync & Async When Consumers Are Behind Firewalls

Since consumers can't receive push notifications or connect to your message broker, a request-acknowledge-poll pattern is the most reliable fallback. Here's the step-by-step flow:

  1. Consumer initiates request: The consumer sends a standard REST POST request to your Flask service with task data.
  2. Service returns 202 Accepted: Your service generates a unique request_id (your GUID), immediately returns a 202 response with the request_id, and drops the task into your message broker queue.
  3. Worker processes task: A dedicated worker (configured to run serially) pulls tasks from the queue, executes the slow processing, and saves the result to a centralized registry tied to the request_id.
  4. Consumer polls for results: The consumer periodically hits a /results/{request_id} endpoint until it receives a completed, failed, or timeout status.

Implementing the Request/Result Registry

You need a fast, accessible store to map request_ids to task statuses and results. Here are two production-ready options:

Option 1: Redis (Lightweight, High-Performance)

Ideal for short-lived results (you can set TTLs to auto-clean old entries):

  • Use request:{request_id} as the key, with a JSON value containing status (pending/processing/completed/failed), result_data, and timestamps.
  • Example Python code with redis-py:
import redis
import json
from datetime import datetime, timedelta

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

def upsert_task_status(request_id: str, status: str, result: dict = None):
    task_data = {
        "status": status,
        "result": result,
        "updated_at": str(datetime.utcnow())
    }
    # Auto-expire after 24 hours to avoid bloat
    redis_client.setex(f"request:{request_id}", timedelta(hours=24), json.dumps(task_data))

def get_task_result(request_id: str) -> dict | None:
    raw_data = redis_client.get(f"request:{request_id}")
    return json.loads(raw_data) if raw_data else None

Option 2: PostgreSQL/MongoDB (Persistent, Long-Term Storage)

Use this if you need to retain results for auditing or compliance:

  • Create a table/collection with fields: request_id (primary key), status, result_data, created_at, updated_at.
  • For Flask, use SQLAlchemy (PostgreSQL) or PyMongo (MongoDB) to interact with the store.

Message Middleware Producer-Consumer Interaction

Let's use RabbitMQ as an example (adjustable for Kafka):

Producer (Flask Service)

When handling the initial consumer request:

import pika
import json
from flask import Flask, request, jsonify

app = Flask(__name__)

@app.route("/submit-task", methods=["POST"])
def submit_task():
    task_data = request.json
    request_id = generate_your_guid()  # Replace with your GUID logic
    
    # Send task to RabbitMQ queue
    connection = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
    channel = connection.channel()
    channel.queue_declare(queue="serial_tasks", durable=True)  # Durable = survive broker restarts
    channel.basic_publish(
        exchange="",
        routing_key="serial_tasks",
        body=json.dumps({"request_id": request_id, "data": task_data}),
        properties=pika.BasicProperties(delivery_mode=pika.DeliveryMode.Persistent)
    )
    connection.close()
    
    # Initialize task status in registry
    upsert_task_status(request_id, "pending")
    
    return jsonify({"request_id": request_id}), 202

Consumer (Serial Worker)

A standalone script (run separately from Gunicorn) to process tasks one at a time:

import pika
import json
from your_registry_module import upsert_task_status

def process_task(ch, method, properties, body):
    task_payload = json.loads(body)
    request_id = task_payload["request_id"]
    task_data = task_payload["data"]
    
    try:
        upsert_task_status(request_id, "processing")
        # Your 30-second slow processing logic here
        result = run_your_slow_worker(task_data)
        upsert_task_status(request_id, "completed", result)
    except Exception as e:
        upsert_task_status(request_id, "failed", str(e))
    finally:
        ch.basic_ack(delivery_tag=method.delivery_tag)  # Acknowledge task completion

def start_serial_worker():
    connection = pika.BlockingConnection(pika.ConnectionParameters("localhost"))
    channel = connection.channel()
    channel.queue_declare(queue="serial_tasks", durable=True)
    # Enforce serial processing: only 1 task at a time per worker
    channel.basic_qos(prefetch_count=1)
    channel.basic_consume(queue="serial_tasks", on_message_callback=process_task)
    
    print("Serial worker started—waiting for tasks...")
    channel.start_consuming()

if __name__ == "__main__":
    start_serial_worker()

Do You Have to Poll? Alternatives (With Caveats)

Polling is the safest bet for firewall-restricted consumers, but you can optimize it:

  • Long Polling: Instead of the consumer making frequent short requests, your /results/{request_id} endpoint holds the connection open until the task completes or hits a timeout (e.g., 30 seconds). This reduces unnecessary HTTP traffic. Example Flask implementation:
from flask import Response, stream_with_context
import time

@app.route("/results/<request_id>")
def long_poll_result(request_id):
    def generate():
        start_time = time.time()
        while time.time() - start_time < 30:  # 30-second timeout
            result = get_task_result(request_id)
            if result and result["status"] in ["completed", "failed"]:
                yield json.dumps(result) + "\n"
                break
            time.sleep(2)  # Check every 2 seconds
        yield json.dumps({"status": "pending"}) + "\n"
    
    return Response(stream_with_context(generate()), mimetype="application/json")
  • Reverse Proxy Workaround: If the consumer can set up a temporary reverse proxy (e.g., ngrok), you could use webhooks to push results—but this depends on the consumer's ability to expose an endpoint, which you noted might not be feasible.

Adapting Your Flask/Gunicorn Setup

  • Your existing Gunicorn sync workers are fine for the Flask producer—sending messages to RabbitMQ is a fast, non-blocking operation.
  • Run your serial worker as a separate process (not under Gunicorn) to avoid blocking web requests. For easier management, use a process manager like systemd or supervisord to keep the worker running.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.13 08:11:53