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 471d35b  Use single instances for compression codec providers (#2017)
471d35b is described below

commit 471d35ba4cf7bc7ed60e82a58c91049a7fbb0262
Author: Matteo Merli <[email protected]>
AuthorDate: Thu Jun 21 19:05:49 2018 -0700

    Use single instances for compression codec providers (#2017)
---
 .../apache/pulsar/client/impl/ConsumerImpl.java    |  4 +-
 .../compression/CompressionCodecProvider.java      | 41 ++++---------------
 .../common/compression/CompressionCodecZLib.java   | 46 +++++++++++++---------
 .../common/compression/CompressorCodecTest.java    |  6 +--
 4 files changed, 39 insertions(+), 58 deletions(-)

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 dd3b5c4..d94fe8b 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 @@ public class ConsumerImpl<T> extends ConsumerBase<T> 
implements ConnectionHandle
     private final int partitionIndex;
 
     private final int receiverQueueRefillThreshold;
-    private final CompressionCodecProvider codecProvider;
 
     private volatile boolean waitingOnReceiveForZeroQueueSize = false;
 
@@ -149,7 +148,6 @@ public class ConsumerImpl<T> extends ConsumerBase<T> 
implements ConnectionHandle
         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 @@ public class ConsumerImpl<T> extends ConsumerBase<T> 
implements ConnectionHandle
     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 212412a..6de8809 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 @@ package org.apache.pulsar.common.compression;
 
 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 bb2aa46..169e347 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 java.util.zip.Inflater;
 
 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 class CompressionCodecZLib implements 
CompressionCodec {
             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 class CompressionCodecZLib implements 
CompressionCodec {
         }
 
         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 654623d..ffa2120 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 @@ public class CompressorCodecTest {
 
     @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);

Reply via email to