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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 04:20:23