基于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()方法(若使用)会一直阻塞
运行步骤
- 安装依赖:执行
pip install flask requests - 先启动消费者服务器:运行消费者代码,等待监听5001端口
- 再启动生产者服务器:运行生产者代码,启动线程发送数据
- 查看控制台输出,可观察数据从生成→发送→接收→处理→回传的完整流程
扩展方向
- 添加失败重试机制:消费者回传失败时将数据重新放入队列
- 替换队列实现:用Redis队列替代内存队列,实现跨进程/服务器的持久化存储
- 引入日志系统:用
logging模块替代print,便于生产环境排查问题
内容的提问来源于stack exchange,提问作者Zurab
相关产品推荐
相关产品推荐

