diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ProducerConsumerInternalTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ProducerConsumerInternalTest.java index 4617f631207f7..503afbefb78c9 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ProducerConsumerInternalTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/impl/ProducerConsumerInternalTest.java @@ -237,7 +237,7 @@ public void testRetentionPolicyByProducingMessages() throws Exception { @Test public void testProducerCompressionMinMsgBodySize() throws PulsarClientException { - byte[] msg1022 = new byte[1022]; + byte[] msg1024 = new byte[1024]; byte[] msg1025 = new byte[1025]; final String topicName = BrokerTestUtil.newUniqueName("persistent://my-property/my-ns/tp_"); @Cleanup @@ -256,7 +256,7 @@ public void testProducerCompressionMinMsgBodySize() throws PulsarClientException producer.conf.setCompressionType(CompressionType.LZ4); // disable batch producer.conf.setBatchingEnabled(false); - producer.newMessage().value(msg1022).send(); + producer.newMessage().value(msg1024).send(); MessageImpl message = (MessageImpl) consumer.receive(); CompressionType compressionType = message.getCompressionType(); assertEquals(compressionType, CompressionType.NONE); @@ -267,7 +267,7 @@ public void testProducerCompressionMinMsgBodySize() throws PulsarClientException // enable batch producer.conf.setBatchingEnabled(true); - producer.newMessage().value(msg1022).send(); + producer.newMessage().value(msg1024).send(); message = (MessageImpl) consumer.receive(); compressionType = message.getCompressionType(); assertEquals(compressionType, CompressionType.NONE); @@ -275,5 +275,20 @@ public void testProducerCompressionMinMsgBodySize() throws PulsarClientException message = (MessageImpl) consumer.receive(); compressionType = message.getCompressionType(); assertEquals(compressionType, CompressionType.LZ4); + + // Verify data integrity + String data = "compression test message"; + producer.conf.setBatchingEnabled(true); + producer.getConfiguration().setCompressMinMsgBodySize(1); + producer.newMessage().value(data.getBytes()).send(); + message = (MessageImpl) consumer.receive(); + assertEquals(new String(message.getData()), data); + + producer.conf.setBatchingEnabled(false); + producer.getConfiguration().setCompressMinMsgBodySize(1); + producer.newMessage().value(data.getBytes()).send(); + message = (MessageImpl) consumer.receive(); + assertEquals(new String(message.getData()), data); + } } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 2752dce945f6b..41df75d9b3e3b 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -542,10 +542,8 @@ public void sendAsync(Message message, SendCallback callback) { boolean compressed = false; // Batch will be compressed when closed // If a message has a delayed delivery time, we'll always send it individually - if (((!isBatchMessagingEnabled() || msgMetadata.hasDeliverAtTime()))) { - if (payload.readableBytes() < conf.getCompressMinMsgBodySize()) { - - } else { + if (!isBatchMessagingEnabled() || msgMetadata.hasDeliverAtTime()) { + if (payload.readableBytes() > conf.getCompressMinMsgBodySize()) { compressedPayload = applyCompression(payload); compressed = true;