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