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

如何创建模拟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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 08:50:45