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

基于Python的多线程HTTP服务器间生产者-消费者双向通信方案问询

双HTTP服务器多线程通信实现(生产者-消费者模型)

1. 生产者服务器(Producer Server)

负责启动6个线程生成随机数据,通过POST请求发送给消费者服务器,同时提供接口接收消费者回传的数据。

from flask import Flask, request
import threading
import requests
import time
import random
import uuid

app = Flask(__name__)

# 存储消费者回传的数据
received_responses = []

def produce_data(consumer_url):
    """单个生产者线程任务:生成数据并发送给消费者"""
    while True:
        # 生成带唯一标识的模拟数据
        data = {
            "message_id": str(uuid.uuid4()),
            "content": f"Data from thread {threading.get_ident()}",
            "timestamp": time.time()
        }
        try:
            # 发送POST请求到消费者接收接口
            response = requests.post(f"{consumer_url}/receive", json=data, timeout=5)
            if response.status_code == 200:
                print(f"Thread {threading.get_ident()}: Sent {data['message_id']} successfully")
            else:
                print(f"Thread {threading.get_ident()}: Failed to send {data['message_id']} - Code {response.status_code}")
        except Exception as e:
            print(f"Thread {threading.get_ident()}: Error sending data - {str(e)}")
        # 随机间隔避免请求过载
        time.sleep(random.uniform(0.5, 2))

@app.route("/feedback", methods=["POST"])
def handle_feedback():
    """接收消费者回传的处理结果"""
    feedback_data = request.get_json()
    received_responses.append(feedback_data)
    print(f"Got feedback: {feedback_data['message_id']} - {feedback_data['status']}")
    return {"status": "ok"}, 200

if __name__ == "__main__":
    CONSUMER_URL = "http://localhost:5001"
    # 启动6个生产者线程
    for _ in range(6):
        thread = threading.Thread(target=produce_data, args=(CONSUMER_URL,), daemon=True)
        thread.start()
    # 启动服务器,监听5000端口
    app.run(host="0.0.0.0", port=5000, threaded=True)

关键细节

  • 用daemon=True设置线程为守护线程,主进程退出时自动终止
  • 每个线程生成带UUID的唯一数据,便于追踪流转过程
  • /feedback接口接收消费者的处理反馈,存入本地列表

2. 消费者服务器(Consumer Server)

负责接收生产者数据并存入线程安全队列,同时启动多线程从队列取数据,处理后回传给生产者。

from flask import Flask, request
import threading
import requests
import time
from queue import Queue

app = Flask(__name__)

# 线程安全队列,存储待处理的生产者数据
data_queue = Queue(maxsize=100)

def process_and_feedback(producer_url):
    """单个消费者线程任务:处理队列数据并回传结果"""
    while True:
        # 阻塞等待队列中的数据
        data = data_queue.get()
        try:
            # 模拟业务处理逻辑
            processed_data = {
                "message_id": data["message_id"],
                "status": "processed_success",
                "processed_time": time.time()
            }
            # 回传处理结果到生产者
            response = requests.post(f"{producer_url}/feedback", json=processed_data, timeout=5)
            if response.status_code == 200:
                print(f"Processed {data['message_id']} and sent feedback")
            else:
                print(f"Failed to send feedback for {data['message_id']} - Code {response.status_code}")
        except Exception as e:
            print(f"Error processing {data['message_id']} - {str(e)}")
        finally:
            # 标记队列任务完成,避免队列阻塞
            data_queue.task_done()

@app.route("/receive", methods=["POST"])
def receive_data():
    """接收生产者发送的数据并放入队列"""
    incoming_data = request.get_json()
    if incoming_data:
        data_queue.put(incoming_data)
        print(f"Received {incoming_data['message_id']}, queue size: {data_queue.qsize()}")
        return {"status": "received"}, 200
    return {"status": "invalid_data"}, 400

if __name__ == "__main__":
    PRODUCER_URL = "http://localhost:5000"
    # 启动4个消费者处理线程
    for _ in range(4):
        thread = threading.Thread(target=process_and_feedback, args=(PRODUCER_URL,), daemon=True)
        thread.start()
    # 启动服务器,监听5001端口
    app.run(host="0.0.0.0", port=5001, threaded=True)

关键细节

  • 使用queue.Queue实现线程安全的数据共享,无需手动加锁
  • data_queue.get()会阻塞直到有数据可用,避免空轮询浪费资源
  • task_done()必须调用,否则队列的join()方法(若使用)会一直阻塞

运行步骤

  1. 安装依赖:执行pip install flask requests
  2. 先启动消费者服务器:运行消费者代码,等待监听5001端口
  3. 再启动生产者服务器:运行生产者代码,启动线程发送数据
  4. 查看控制台输出,可观察数据从生成→发送→接收→处理→回传的完整流程

扩展方向

  • 添加失败重试机制:消费者回传失败时将数据重新放入队列
  • 替换队列实现:用Redis队列替代内存队列,实现跨进程/服务器的持久化存储
  • 引入日志系统:用logging模块替代print,便于生产环境排查问题

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 04:05:21