如何创建模拟Kafka生产者与消费者?无Kafka环境实现方案咨询
模拟Kafka生产者-消费者的实现思路与代码示例
完全可以用Python实现轻量的模拟版本,无需安装Kafka,核心是用线程安全队列模拟消息中间件的存储转发逻辑:生产者从JSON文件读取消息并推入队列,消费者从队列拉取消息处理。
实现步骤与代码
1. 准备JSON消息文件
先创建messages.json文件,格式为JSON数组,每条消息为一个JSON对象:
[ {"id": 1, "content": "第一条测试消息"}, {"id": 2, "content": "第二条测试消息"}, {"id": 3, "content": "第三条测试消息"} ]
2. 模拟生产者与消费者代码
import json import queue import threading import time # 全局队列,模拟Kafka的Topic存储 message_queue = queue.Queue() # 标记生产者是否完成任务 producer_done = False def producer(file_path): global producer_done try: with open(file_path, 'r', encoding='utf-8') as f: messages = json.load(f) for msg in messages: # 模拟Kafka生产者发送延迟 time.sleep(0.5) message_queue.put(json.dumps(msg)) print(f"生产者已发送消息: {json.dumps(msg)}") except Exception as e: print(f"生产者出错: {str(e)}") finally: producer_done = True print("生产者已完成所有消息发送") def consumer(): while True: # 队列空且生产者完成时退出 if message_queue.empty() and producer_done: print("消费者已处理完所有消息,退出") break try: # 阻塞等待消息,超时1秒避免死循环 msg = message_queue.get(timeout=1) # 模拟业务处理逻辑,可替换为你的代码 parsed_msg = json.loads(msg) print(f"消费者已处理消息: {parsed_msg}") # 标记消息处理完成(类似Kafka的Offset提交) message_queue.task_done() except queue.Empty: continue if __name__ == "__main__": # 启动生产者线程 prod_thread = threading.Thread(target=producer, args=("messages.json",)) prod_thread.start() # 启动消费者线程(可启动多个,模拟多消费者组) cons_thread = threading.Thread(target=consumer) cons_thread.start() # 等待生产者完成 prod_thread.join() # 等待队列所有消息处理完毕 message_queue.join() # 等待消费者退出 cons_thread.join()
关键说明
- 依赖:仅用Python标准库,无需额外安装包,他人拿到代码可直接运行
- 扩展性:可轻松修改为多消费者模式(启动多个consumer线程),或添加消息重试、本地持久化等逻辑
- 适配性:若你的JSON文件是每行一个JSON字符串(非数组格式),只需修改生产者读取逻辑为逐行读取即可
内容的提问来源于stack exchange,提问作者Oxana Grey
相关产品推荐
相关产品推荐

