From 5a602a380a3772dbcc409e3547e9a68bd1370bc7 Mon Sep 17 00:00:00 2001 From: xiangying Date: Thu, 20 Mar 2025 21:37:52 +0800 Subject: [PATCH 1/3] [fix][client] Unify the data size check during compression ## Motivation Unify the data size check during compression. Currently, there is inconsistency in the criteria for determining the threshold of messages with compression enabled: For instance, with compressMinMsgBodySize set to 4kb: - When batching is not enabled, compression is applied to messages that are equal to or greater than 4kb. - When batching is enabled, compression is applied to messages that are greater than 4kb. ## Modifications 1. Standardize the criteria for determining the threshold to enable compression. 2. Optimize formatting. 3. Enhance testing. --- .../client/impl/ProducerConsumerInternalTest.java | 15 +++++++++++++++ .../apache/pulsar/client/impl/ProducerImpl.java | 6 ++---- 2 files changed, 17 insertions(+), 4 deletions(-) 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..0a9fb894c8459 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 @@ -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..4a8f77b822696 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; From 098c863bf541f79114858e48bf06d5fdf48fc803 Mon Sep 17 00:00:00 2001 From: xiangying Date: Mon, 24 Mar 2025 14:47:24 +0800 Subject: [PATCH 2/3] fix --- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 4a8f77b822696..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,7 +542,7 @@ 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 (!isBatchMessagingEnabled() || msgMetadata.hasDeliverAtTime()) { if (payload.readableBytes() > conf.getCompressMinMsgBodySize()) { compressedPayload = applyCompression(payload); compressed = true; From a765f5c3172415a43156a4387f9ac15e930aef43 Mon Sep 17 00:00:00 2001 From: xiangying Date: Tue, 25 Mar 2025 17:44:04 +0800 Subject: [PATCH 3/3] [test] Covers the case where the message size is equal to the threshold --- .../pulsar/client/impl/ProducerConsumerInternalTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 0a9fb894c8459..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);