From 41935c97929846a9082122ed73c25d484ddbc302 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=B8=9B=E6=90=8F?= <237716289@qq.com> Date: Tue, 30 Jul 2019 12:04:53 +0800 Subject: [PATCH] Consumer can getValue return null --- .../apache/pulsar/client/impl/MessageImpl.java | 6 +++++- .../pulsar/client/impl/MessageImplTest.java | 16 ++++++++++++++++ 2 files changed, 21 insertions(+), 1 deletion(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java index 3398d36ed22fe..b47e0e37f1074 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/MessageImpl.java @@ -270,7 +270,11 @@ public T getValue() { return schema.decode(getData(), schemaVersion); } } else { - return schema.decode(getData()); + if (getData().length == 0) { + return null; + } else { + return schema.decode(getData()); + } } } } diff --git a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java index dd4944512841a..ebb33d6a5cdd8 100644 --- a/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java +++ b/pulsar-client/src/test/java/org/apache/pulsar/client/impl/MessageImplTest.java @@ -25,6 +25,7 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.schema.SchemaDefinition; import org.apache.pulsar.client.impl.schema.AvroSchema; +import org.apache.pulsar.client.impl.schema.BooleanSchema; import org.apache.pulsar.client.impl.schema.JSONSchema; import org.apache.pulsar.client.impl.schema.SchemaTestUtils; import org.apache.pulsar.client.impl.schema.generic.MultiVersionSchemaInfoProvider; @@ -398,4 +399,19 @@ public void testSeparatedAVROJSONVersionGetProducerDataAssigned() { KeyValueEncodingType.valueOf(keyValueSchema.getSchemaInfo().getProperties().get("kv.encoding.type")), KeyValueEncodingType.SEPARATED); } + + @Test + public void testTypedSchemaGetNullValue() { + + byte[] encodeBytes = new byte[0]; + MessageMetadata.Builder builder = MessageMetadata.newBuilder() + .setProducerName("getNullValue"); + ByteString byteString = ByteString.copyFrom(new byte[0]); + builder.setSchemaVersion(byteString); + builder.setPartitionKey(Base64.getEncoder().encodeToString(encodeBytes)); + builder.setPartitionKeyB64Encoded(true); + MessageImpl msg = MessageImpl.create( + builder, ByteBuffer.wrap(encodeBytes), BooleanSchema.of()); + assertNull(msg.getValue()); + } }