This is an automated email from the ASF dual-hosted git repository.
mmerli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new 44cbed3 Encrypted or compressed payloads can be compacted out if
empty (#1889)
44cbed3 is described below
commit 44cbed3a842f9a04e656b987433f686e56c616a9
Author: Ivan Kelly <[email protected]>
AuthorDate: Mon Jun 4 19:08:18 2018 +0200
Encrypted or compressed payloads can be compacted out if empty (#1889)
Compressed payloads don't have zero size, even if the uncompressed
payload is empty, so if a payload was empty, it wouldn't delete that key
from the compaction result.
Similarly for encrypted messages, zero size unencrypted doesn't map to
zero size encrypted.
This patch adds special handling to use getUncompressedSize() to find
the size of the payload, rather than looking at the readableBytes of
the payload. UncompressedSize is now set for all messages.
We still fall back to using readableBytes for the case where we have
an old client that isn't updating UncompressedSize.
---
.../pulsar/compaction/TwoPhaseCompactor.java | 7 +-
.../apache/pulsar/compaction/CompactionTest.java | 157 +++++++++++++++++++++
.../apache/pulsar/client/impl/ProducerImpl.java | 2 +-
3 files changed, 164 insertions(+), 2 deletions(-)
diff --git
a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java
b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java
index 5db4f2e..b4ee68a 100644
---
a/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java
+++
b/pulsar-broker/src/main/java/org/apache/pulsar/compaction/TwoPhaseCompactor.java
@@ -47,6 +47,7 @@ import org.apache.pulsar.client.api.RawMessage;
import org.apache.pulsar.client.api.RawReader;
import org.apache.pulsar.client.impl.RawBatchConverter;
import org.apache.pulsar.common.api.Commands;
+import org.apache.pulsar.common.api.proto.PulsarApi.CompressionType;
import org.apache.pulsar.common.api.proto.PulsarApi.MessageMetadata;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -314,7 +315,11 @@ public class TwoPhaseCompactor extends Compactor {
MessageMetadata msgMetadata =
Commands.parseMessageMetadata(headersAndPayload);
try {
if (msgMetadata.hasPartitionKey()) {
- return Pair.of(msgMetadata.getPartitionKey(),
headersAndPayload.readableBytes());
+ int size = headersAndPayload.readableBytes();
+ if (msgMetadata.hasUncompressedSize()) {
+ size = msgMetadata.getUncompressedSize();
+ }
+ return Pair.of(msgMetadata.getPartitionKey(), size);
} else {
return null;
}
diff --git
a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactionTest.java
b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactionTest.java
index d980465..95a10fa 100644
---
a/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactionTest.java
+++
b/pulsar-broker/src/test/java/org/apache/pulsar/compaction/CompactionTest.java
@@ -664,6 +664,83 @@ public class CompactionTest extends
MockedPulsarServiceBaseTest {
}
}
+ @Test
+ public void testEmptyPayloadDeletesWhenCompressed() throws Exception {
+ String topic = "persistent://my-property/use/my-ns/my-topic1";
+
+ // subscribe before sending anything, so that we get all messages
+ pulsarClient.newConsumer().topic(topic).subscriptionName("sub1")
+ .readCompacted(true).subscribe().close();
+
+ try (Producer<byte[]> producerNormal = pulsarClient.newProducer()
+ .topic(topic)
+ .enableBatching(false)
+ .compressionType(CompressionType.LZ4)
+ .create();
+ Producer<byte[]> producerBatch = pulsarClient.newProducer()
+ .topic(topic)
+ .maxPendingMessages(3)
+ .enableBatching(true)
+ .compressionType(CompressionType.LZ4)
+ .batchingMaxMessages(3)
+ .batchingMaxPublishDelay(1, TimeUnit.HOURS)
+ .create()) {
+
+ // key0 persists through it all
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key0")
+
.setContent("my-message-0".getBytes()).build()).get();
+
+ // key1 is added but then deleted
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key1")
+
.setContent("my-message-1".getBytes()).build()).get();
+
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key1").build()).get();
+
+ // key2 is added but deleted in same batch
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key2")
+
.setContent("my-message-2".getBytes()).build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key3")
+
.setContent("my-message-3".getBytes()).build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key2").build()).get();
+
+ // key3 is added in previous batch, deleted in this batch
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key3").build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key4")
+
.setContent("my-message-3".getBytes()).build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key4").build()).get();
+
+ // key4 is added, deleted, then resurrected
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key4")
+
.setContent("my-message-4".getBytes()).build()).get();
+ }
+
+ // compact the topic
+ Compactor compactor = new TwoPhaseCompactor(conf, pulsarClient, bk,
compactionScheduler);
+ compactor.compact(topic).get();
+
+ try (Consumer<byte[]> consumer =
pulsarClient.newConsumer().topic(topic)
+ .subscriptionName("sub1").readCompacted(true).subscribe()){
+ Message<byte[]> message1 = consumer.receive();
+ Assert.assertEquals(message1.getKey(), "key0");
+ Assert.assertEquals(new String(message1.getData()),
"my-message-0");
+
+ Message<byte[]> message2 = consumer.receive();
+ Assert.assertEquals(message2.getKey(), "key4");
+ Assert.assertEquals(new String(message2.getData()),
"my-message-4");
+ }
+ }
+
+ // test compact no keys
@Test
public void testCompactorReadsCompacted() throws Exception {
@@ -1055,4 +1132,84 @@ public class CompactionTest extends
MockedPulsarServiceBaseTest {
}
}
+ @Test
+ public void testEmptyPayloadDeletesWhenEncrypted() throws Exception {
+ String topic = "persistent://my-property/use/my-ns/my-topic1";
+
+ // subscribe before sending anything, so that we get all messages
+ pulsarClient.newConsumer().topic(topic).subscriptionName("sub1")
+ .readCompacted(true).subscribe().close();
+
+ try (Producer<byte[]> producerNormal = pulsarClient.newProducer()
+ .topic(topic)
+ .enableBatching(false)
+ .addEncryptionKey("client-ecdsa.pem").cryptoKeyReader(new
EncKeyReader())
+ .create();
+ Producer<byte[]> producerBatch = pulsarClient.newProducer()
+ .topic(topic)
+ .maxPendingMessages(3)
+ .enableBatching(true)
+ .addEncryptionKey("client-ecdsa.pem").cryptoKeyReader(new
EncKeyReader())
+ .batchingMaxMessages(3)
+ .batchingMaxPublishDelay(1, TimeUnit.HOURS)
+ .create()) {
+
+ // key0 persists through it all
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key0")
+
.setContent("my-message-0".getBytes()).build()).get();
+
+ // key1 is added but then deleted
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key1")
+
.setContent("my-message-1".getBytes()).build()).get();
+
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key1").build()).get();
+
+ // key2 is added but deleted in same batch
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key2")
+
.setContent("my-message-2".getBytes()).build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key3")
+
.setContent("my-message-3".getBytes()).build());
+ producerBatch.sendAsync(MessageBuilder.create()
+ .setKey("key2").build()).get();
+
+ // key4 is added, deleted, then resurrected
+ producerNormal.sendAsync(MessageBuilder.create()
+ .setKey("key4")
+
.setContent("my-message-4".getBytes()).build()).get();
+ }
+
+ // compact the topic
+ Compactor compactor = new TwoPhaseCompactor(conf, pulsarClient, bk,
compactionScheduler);
+ compactor.compact(topic).get();
+
+ try (Consumer<byte[]> consumer =
pulsarClient.newConsumer().topic(topic)
+ .cryptoKeyReader(new EncKeyReader())
+ .subscriptionName("sub1").readCompacted(true).subscribe()){
+ Message<byte[]> message1 = consumer.receive();
+ Assert.assertEquals(message1.getKey(), "key0");
+ Assert.assertEquals(new String(message1.getData()),
"my-message-0");
+
+ // see all messages from batch
+ Message<byte[]> message2 = consumer.receive();
+ Assert.assertEquals(message2.getKey(), "key2");
+ Assert.assertEquals(new String(message2.getData()),
"my-message-2");
+
+ Message<byte[]> message3 = consumer.receive();
+ Assert.assertEquals(message3.getKey(), "key3");
+ Assert.assertEquals(new String(message3.getData()),
"my-message-3");
+
+ Message<byte[]> message4 = consumer.receive();
+ Assert.assertEquals(message4.getKey(), "key2");
+ Assert.assertEquals(new String(message4.getData()), "");
+
+ Message<byte[]> message5 = consumer.receive();
+ Assert.assertEquals(message5.getKey(), "key4");
+ Assert.assertEquals(new String(message5.getData()),
"my-message-4");
+ }
+ }
}
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 875fb87..b5962ac 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
@@ -313,8 +313,8 @@ public class ProducerImpl<T> extends ProducerBase<T>
implements TimerTask, Conne
if (conf.getCompressionType() != CompressionType.NONE) {
msgMetadataBuilder.setCompression(convertCompressionType(conf.getCompressionType()));
-
msgMetadataBuilder.setUncompressedSize(uncompressedSize);
}
+ msgMetadataBuilder.setUncompressedSize(uncompressedSize);
}
if (isBatchMessagingEnabled()) {
--
To stop receiving notification emails like this one, please contact
[email protected].