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);