mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-30 18:10:44 +08:00
Accomodate updated openmessaging api
This commit is contained in:
@@ -20,6 +20,7 @@ import io.openmessaging.BytesMessage;
|
||||
import io.openmessaging.KeyValue;
|
||||
import io.openmessaging.Message;
|
||||
import io.openmessaging.OMS;
|
||||
import io.openmessaging.exception.OMSMessageFormatException;
|
||||
import org.apache.commons.lang3.builder.ToStringBuilder;
|
||||
|
||||
public class BytesMessageImpl implements BytesMessage {
|
||||
@@ -33,8 +34,12 @@ public class BytesMessageImpl implements BytesMessage {
|
||||
}
|
||||
|
||||
@Override
|
||||
public byte[] getBody() {
|
||||
return body;
|
||||
public <T> T getBody(Class<T> type) throws OMSMessageFormatException {
|
||||
if (type == byte[].class) {
|
||||
return (T)body;
|
||||
}
|
||||
|
||||
throw new OMSMessageFormatException("", "Cannot assign byte[] to " + type.getName());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -46,7 +46,7 @@ public class OMSUtil {
|
||||
|
||||
public static org.apache.rocketmq.common.message.Message msgConvert(BytesMessage omsMessage) {
|
||||
org.apache.rocketmq.common.message.Message rmqMessage = new org.apache.rocketmq.common.message.Message();
|
||||
rmqMessage.setBody(omsMessage.getBody());
|
||||
rmqMessage.setBody(omsMessage.getBody(byte[].class));
|
||||
|
||||
KeyValue sysHeaders = omsMessage.sysHeaders();
|
||||
KeyValue userHeaders = omsMessage.userHeaders();
|
||||
|
||||
+1
-1
@@ -83,7 +83,7 @@ public class PullConsumerImplTest {
|
||||
|
||||
Message message = consumer.receive();
|
||||
assertThat(message.sysHeaders().getString(Message.BuiltinKeys.MESSAGE_ID)).isEqualTo("NewMsgId");
|
||||
assertThat(((BytesMessage) message).getBody()).isEqualTo(testBody);
|
||||
assertThat(((BytesMessage) message).getBody(byte[].class)).isEqualTo(testBody);
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
+1
-1
@@ -75,7 +75,7 @@ public class PushConsumerImplTest {
|
||||
@Override
|
||||
public void onReceived(Message message, Context context) {
|
||||
assertThat(message.sysHeaders().getString(Message.BuiltinKeys.MESSAGE_ID)).isEqualTo("NewMsgId");
|
||||
assertThat(((BytesMessage) message).getBody()).isEqualTo(testBody);
|
||||
assertThat(((BytesMessage) message).getBody(byte[].class)).isEqualTo(testBody);
|
||||
context.ack();
|
||||
}
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user