如何用Pika、Python和RabbitMQ发送GZIP压缩JSON消息
使用Pika、Python和RabbitMQ发送GZIP压缩的大JSON消息
核心方案
将大JSON序列化为字节串后用GZIP压缩,发送压缩后的二进制数据;接收端识别压缩格式后解压并还原为JSON。这种方式能大幅降低消息体积,轻松绕过RabbitMQ的128MB大小限制。
发送端代码(压缩并发送)
import pika import json import gzip # 替换为你的实际大JSON数据 large_json_data = {"sample_key": "sample_value" * 10**8} # 模拟500MB级别的JSON数据 # 1. 把JSON序列化为UTF-8编码的字节串 json_bytes = json.dumps(large_json_data).encode("utf-8") # 2. GZIP压缩字节串 compressed_bytes = gzip.compress(json_bytes) print(f"原始大小: {len(json_bytes)/1024/1024:.2f} MB | 压缩后大小: {len(compressed_bytes)/1024/1024:.2f} MB") # 3. 连接RabbitMQ并发送压缩消息 conn = pika.BlockingConnection(pika.ConnectionParameters("localhost")) channel = conn.channel() # 声明队列(如需持久化,添加durable=True) channel.queue_declare(queue="large_compressed_queue") # 发送消息,通过content_type标记为GZIP格式 channel.basic_publish( exchange="", routing_key="large_compressed_queue", body=compressed_bytes, properties=pika.BasicProperties( content_type="application/gzip", # 告知接收端消息是GZIP压缩的 delivery_mode=pika.DeliveryMode.Persistent # 可选:持久化消息防止RabbitMQ重启丢失 ) ) conn.close()
接收端代码(解压并解析)
import pika import json import gzip def process_compressed_message(ch, method, properties, body): # 验证消息是否为GZIP压缩格式 if properties.content_type == "application/gzip": try: # 1. 解压二进制数据 decompressed_bytes = gzip.decompress(body) # 2. 还原为JSON对象 json_data = json.loads(decompressed_bytes.decode("utf-8")) print(f"成功解析JSON,示例字段长度: {len(json_data['sample_key'])}") except Exception as e: print(f"处理消息失败: {str(e)}") else: print("忽略非GZIP格式的消息") # 手动确认消息已处理完成 ch.basic_ack(delivery_tag=method.delivery_tag) # 建立RabbitMQ连接 conn = pika.BlockingConnection(pika.ConnectionParameters("localhost")) channel = conn.channel() channel.queue_declare(queue="large_compressed_queue") # 设置预取计数,避免一次性接收过多消息导致内存溢出 channel.basic_qos(prefetch_count=1) # 绑定消息处理函数 channel.basic_consume(queue="large_compressed_queue", on_message_callback=process_compressed_message) print("等待接收压缩消息...") channel.start_consuming()
关键注意事项
- 压缩率验证:JSON文本的GZIP压缩率通常可达5-10倍,500MB的JSON压缩后一般能控制在100MB以内,可通过代码中的打印值确认
- 持久化配置:若需要消息不丢失,需同时将队列声明为
durable=True,并设置消息的delivery_mode=Persistent - 内存控制:消费者端必须设置
basic_qos,避免一次性拉取大量压缩消息,导致内存占用过高 - 异常处理:实际生产环境中需添加更多异常捕获(如RabbitMQ连接中断、压缩/解压失败、JSON解析错误等)
- 极端情况处理:如果压缩后仍接近128MB,可考虑将数据分块发送,但GZIP压缩已能覆盖绝大多数大JSON场景
内容的提问来源于stack exchange,提问作者GuneetSingh
相关产品推荐
相关产品推荐

