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

如何在EmbeddedActiveMQ消息中发送对象?迁移遇序列化异常

EmbeddedJMS迁移至EmbeddedActiveMQ的对象消息序列化问题

这是从EmbeddedJMS迁移到EmbeddedActiveMQ时遇到的问题。

此前使用JMS发送可序列化对象消息的代码:

try {
    producer.send(session.createObjectMessage(message));
    session.commit();
} catch (JMSException e) {
    log.error(e);
}

改用EmbeddedActiveMQ后,使用以下代码生产消息:

try {
    final ClientMessage msg = clientSession.createMessage(true);
    try(ByteArrayOutputStream baos = new ByteArrayOutputStream();
            ObjectOutputStream oos = new ObjectOutputStream(baos)) {
        oos.writeObject(message);
        oos.flush();
        msg.writeBodyBufferBytes(baos.toByteArray());
        clientProducer.send(msg);
        clientSession.commit();
    }
} catch(ActiveMQException | IOException e) {
    throw new RuntimeException(e);
}

接收消息时的代码:

ByteArrayInputStream arrayInputStream = new ByteArrayInputStream(clientMessage.getBodyBuffer().toByteBuffer().array());
ObjectInputStream objectInputStream = null;
try {
    objectInputStream = new ObjectInputStream(arrayInputStream);
    objectInputStream.readObject();
} catch(IOException | ClassNotFoundException e) {
    throw new RuntimeException(e);
}

运行时抛出异常:

java.io.StreamCorruptedException: invalid stream header: 000000B2

问题原因与解决方案

1. 异常根源

ClientMessage.writeBodyBufferBytes()方法会自动在写入的字节数组前添加4字节的int类型长度前缀(000000B2就是序列化后对象的字节长度十进制178的十六进制表示),但接收时直接把整个ByteBuffer转成字节数组,包含了这4个长度字节,导致ObjectInputStream将长度前缀误判为对象流的头部,从而抛出异常。

2. 修复后的手动序列化实现

生产消息

改用getBodyBuffer().writeBytes()方法直接写入序列化后的字节,避免自动添加长度前缀:

try {
    final ClientMessage msg = clientSession.createMessage(true);
    try(ByteArrayOutputStream baos = new ByteArrayOutputStream();
            ObjectOutputStream oos = new ObjectOutputStream(baos)) {
        oos.writeObject(message);
        oos.flush();
        byte[] serializedData = baos.toByteArray();
        msg.getBodyBuffer().writeBytes(serializedData);
        clientProducer.send(msg);
        clientSession.commit();
    }
} catch(ActiveMQException | IOException e) {
    throw new RuntimeException(e);
}

接收消息

直接读取整个body buffer的字节即可:

try {
    ByteBuffer buffer = clientMessage.getBodyBuffer().toByteBuffer();
    byte[] data = new byte[buffer.remaining()];
    buffer.get(data);
    try(ByteArrayInputStream bais = new ByteArrayInputStream(data);
            ObjectInputStream ois = new ObjectInputStream(bais)) {
        Object targetObj = ois.readObject();
        // 处理反序列化后的对象
    }
} catch(IOException | ClassNotFoundException | ActiveMQException e) {
    throw new RuntimeException(e);
}

如果坚持使用writeBodyBufferBytes(),接收时需要先读取前4字节的长度,再读取对应长度的字节:

try {
    ByteBuffer buffer = clientMessage.getBodyBuffer().toByteBuffer();
    int dataLength = buffer.getInt();
    byte[] data = new byte[dataLength];
    buffer.get(data);
    try(ByteArrayInputStream bais = new ByteArrayInputStream(data);
            ObjectInputStream ois = new ObjectInputStream(bais)) {
        Object targetObj = ois.readObject();
        // 处理反序列化后的对象
    }
} catch(IOException | ClassNotFoundException | ActiveMQException e) {
    throw new RuntimeException(e);
}

3. 更优雅的实现方式

ActiveMQ Artemis(EmbeddedActiveMQ基于它)支持原生的对象消息处理,无需手动序列化:

方式一:使用Core API的属性存储对象

// 生产消息
try {
    ClientMessage msg = clientSession.createMessage(true);
    msg.putObjectProperty("businessObject", message);
    clientProducer.send(msg);
    clientSession.commit();
} catch(ActiveMQException e) {
    throw new RuntimeException(e);
}

// 接收消息
try {
    Object businessObj = clientMessage.getObjectProperty("businessObject");
    // 处理对象
} catch(ActiveMQException e) {
    throw new RuntimeException(e);
}

方式二:兼容JMS API的ObjectMessage

如果希望保留JMS兼容性,直接使用JMS的ObjectMessage即可,无需切换到Core API:

// 生产端
try (JMSContext context = connectionFactory.createContext()) {
    ObjectMessage msg = context.createObjectMessage(message);
    context.createProducer().send(targetQueue, msg);
} catch(JMSException e) {
    throw new RuntimeException(e);
}

// 接收端
try (JMSContext context = connectionFactory.createContext()) {
    ObjectMessage msg = (ObjectMessage) context.createConsumer(targetQueue).receive();
    Serializable businessObj = msg.getObject();
    // 处理对象
} catch(JMSException e) {
    throw new RuntimeException(e);
}

这种方式既符合JMS规范,又避免了手动序列化的繁琐与出错风险,是更推荐的实现方式。


内容的提问来源于stack exchange,提问作者Wojciech Piotrowiak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 01:55:20