merlimat closed pull request #2017: Use single instances for compression codec 
providers
URL: https://github.com/apache/incubator-pulsar/pull/2017
 
 
   

This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:

As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):

diff --git 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
index dd3b5c4675..d94fe8b62d 100644
--- 
a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
+++ 
b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java
@@ -97,7 +97,6 @@
     private final int partitionIndex;
 
     private final int receiverQueueRefillThreshold;
-    private final CompressionCodecProvider codecProvider;
 
     private volatile boolean waitingOnReceiveForZeroQueueSize = false;
 
@@ -149,7 +148,6 @@
         this.subscribeTimeout = System.currentTimeMillis() + 
client.getConfiguration().getOperationTimeoutMs();
         this.partitionIndex = partitionIndex;
         this.receiverQueueRefillThreshold = conf.getReceiverQueueSize() / 2;
-        this.codecProvider = new CompressionCodecProvider();
         this.priorityLevel = conf.getPriorityLevel();
         this.readCompacted = conf.isReadCompacted();
         this.subscriptionInitialPosition = 
conf.getSubscriptionInitialPosition();
@@ -1027,7 +1025,7 @@ private ByteBuf decryptPayloadIfNeeded(MessageIdData 
messageId, MessageMetadata
     private ByteBuf uncompressPayloadIfNeeded(MessageIdData messageId, 
MessageMetadata msgMetadata, ByteBuf payload,
             ClientCnx currentCnx) {
         CompressionType compressionType = msgMetadata.getCompression();
-        CompressionCodec codec = codecProvider.getCodec(compressionType);
+        CompressionCodec codec = 
CompressionCodecProvider.getCompressionCodec(compressionType);
         int uncompressedSize = msgMetadata.getUncompressedSize();
         int payloadSize = payload.readableBytes();
         if (payloadSize > PulsarDecoder.MaxMessageSize) {
diff --git 
a/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecProvider.java
 
b/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecProvider.java
index 212412aabb..6de8809e16 100644
--- 
a/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecProvider.java
+++ 
b/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecProvider.java
@@ -20,47 +20,22 @@
 
 import java.util.EnumMap;
 
+import lombok.experimental.UtilityClass;
+
 import org.apache.pulsar.common.api.proto.PulsarApi.CompressionType;
 
+@UtilityClass
 public class CompressionCodecProvider {
-    private static final EnumMap<CompressionType, Class<? extends 
CompressionCodec>> codecs;
-
-    private final static CompressionCodec compressionCodecNone = new 
CompressionCodecNone();
+    private static final EnumMap<CompressionType, CompressionCodec> codecs;
 
     static {
         codecs = new EnumMap<>(CompressionType.class);
-        codecs.put(CompressionType.NONE, CompressionCodecNone.class);
-        codecs.put(CompressionType.LZ4, CompressionCodecLZ4.class);
-        codecs.put(CompressionType.ZLIB, CompressionCodecZLib.class);
+        codecs.put(CompressionType.NONE, new CompressionCodecNone());
+        codecs.put(CompressionType.LZ4, new CompressionCodecLZ4());
+        codecs.put(CompressionType.ZLIB, new CompressionCodecZLib());
     }
 
     public static CompressionCodec getCompressionCodec(CompressionType type) {
-        if (type == CompressionType.NONE) {
-            // Always use the same instance for the none-codec
-            return compressionCodecNone;
-        }
-
-        try {
-            return codecs.get(type).newInstance();
-        } catch (Exception e) {
-            throw new RuntimeException(e);
-        }
-    }
-
-    private final EnumMap<CompressionType, CompressionCodec> codecInstances;
-
-    public CompressionCodecProvider() {
-        codecInstances = new EnumMap<>(CompressionType.class);
-        try {
-            for (CompressionType type : CompressionType.values()) {
-                codecInstances.put(type, codecs.get(type).newInstance());
-            }
-        } catch (Exception e) {
-            throw new RuntimeException(e);
-        }
-    }
-
-    public CompressionCodec getCodec(CompressionType type) {
-        return codecInstances.get(type);
+        return codecs.get(type);
     }
 }
diff --git 
a/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecZLib.java
 
b/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecZLib.java
index bb2aa4677f..169e34706d 100644
--- 
a/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecZLib.java
+++ 
b/pulsar-common/src/main/java/org/apache/pulsar/common/compression/CompressionCodecZLib.java
@@ -27,14 +27,26 @@
 
 import io.netty.buffer.ByteBuf;
 import io.netty.buffer.PooledByteBufAllocator;
+import io.netty.util.concurrent.FastThreadLocal;
 
 /**
  * ZLib Compression
  */
 public class CompressionCodecZLib implements CompressionCodec {
 
-    private final Deflater deflater = new Deflater();
-    private final Inflater inflater = new Inflater();
+    private final FastThreadLocal<Deflater> deflater = new 
FastThreadLocal<Deflater>() {
+        @Override
+        protected Deflater initialValue() throws Exception {
+            return new Deflater();
+        }
+    };
+
+    private final FastThreadLocal<Inflater> inflater = new 
FastThreadLocal<Inflater>() {
+        @Override
+        protected Inflater initialValue() throws Exception {
+            return new Inflater();
+        }
+    };
 
     @Override
     public ByteBuf encode(ByteBuf source) {
@@ -54,19 +66,17 @@ public ByteBuf encode(ByteBuf source) {
             source.getBytes(source.readerIndex(), array);
         }
 
-        synchronized (deflater) {
-            deflater.setInput(array, offset, length);
-            while (!deflater.needsInput()) {
-                deflate(compressed);
-            }
-
-            deflater.reset();
+        Deflater deflater = this.deflater.get();
+        deflater.reset();
+        deflater.setInput(array, offset, length);
+        while (!deflater.needsInput()) {
+            deflate(deflater, compressed);
         }
 
         return compressed;
     }
 
-    private void deflate(ByteBuf out) {
+    private static void deflate(Deflater deflater, ByteBuf out) {
         int numBytes;
         do {
             int writerIndex = out.writerIndex();
@@ -94,14 +104,14 @@ public ByteBuf decode(ByteBuf encoded, int 
uncompressedLength) throws IOExceptio
         }
 
         int resultLength;
-        synchronized (inflater) {
-            inflater.setInput(array, offset, len);
-            try {
-                resultLength = inflater.inflate(uncompressed.array(), 
uncompressed.arrayOffset(), uncompressedLength);
-            } catch (DataFormatException e) {
-                throw new IOException(e);
-            }
-            inflater.reset();
+        Inflater inflater = this.inflater.get();
+        inflater.reset();
+        inflater.setInput(array, offset, len);
+
+        try {
+            resultLength = inflater.inflate(uncompressed.array(), 
uncompressed.arrayOffset(), uncompressedLength);
+        } catch (DataFormatException e) {
+            throw new IOException(e);
         }
 
         checkArgument(resultLength == uncompressedLength);
diff --git 
a/pulsar-common/src/test/java/org/apache/pulsar/common/compression/CompressorCodecTest.java
 
b/pulsar-common/src/test/java/org/apache/pulsar/common/compression/CompressorCodecTest.java
index 654623d852..ffa21205f0 100644
--- 
a/pulsar-common/src/test/java/org/apache/pulsar/common/compression/CompressorCodecTest.java
+++ 
b/pulsar-common/src/test/java/org/apache/pulsar/common/compression/CompressorCodecTest.java
@@ -109,10 +109,8 @@ void testMultpileUsages(CompressionType type) throws 
IOException {
 
     @Test(dataProvider = "codec")
     void testCodecProvider(CompressionType type) throws IOException {
-        CompressionCodecProvider provider = new CompressionCodecProvider();
-
-        CompressionCodec codec1 = provider.getCodec(type);
-        CompressionCodec codec2 = provider.getCodec(type);
+        CompressionCodec codec1 = 
CompressionCodecProvider.getCompressionCodec(type);
+        CompressionCodec codec2 = 
CompressionCodecProvider.getCompressionCodec(type);
 
         // A single provider instance must return the same codec instance 
every time
         assertTrue(codec1 == codec2);


 

----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
 
For queries about this service, please contact Infrastructure at:
[email protected]


With regards,
Apache Git Services

Reply via email to