From 4c47560c1e5aeaef81bafa19d6c88ee2051f34ca Mon Sep 17 00:00:00 2001 From: Yunze Xu Date: Wed, 2 Feb 2022 21:13:22 +0800 Subject: [PATCH] Fix thread safety issue for MessageFetchContext --- .../handlers/kop/MessageFetchContext.java | 43 ++++++++----------- 1 file changed, 18 insertions(+), 25 deletions(-) diff --git a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java index a79b836d42..165feaad14 100644 --- a/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java +++ b/kafka-impl/src/main/java/io/streamnative/pulsar/handlers/kop/MessageFetchContext.java @@ -83,21 +83,21 @@ protected MessageFetchContext newObject(Handle handle) { }; private final Handle recyclerHandle; - private Map> responseData; - private ConcurrentLinkedQueue decodeResults; - private KafkaRequestHandler requestHandler; - private int maxReadEntriesNum; - private KafkaTopicManager topicManager; - private RequestStats statsLogger; - private TransactionCoordinator tc; - private String clientHost; - private FetchRequest fetchRequest; - private RequestHeader header; + private final Map> responseData = new ConcurrentHashMap<>(); + private final ConcurrentLinkedQueue decodeResults = new ConcurrentLinkedQueue<>(); + private final AtomicBoolean hasComplete = new AtomicBoolean(false); + private final AtomicLong bytesReadable = new AtomicLong(0); + private volatile KafkaRequestHandler requestHandler; + private volatile int maxReadEntriesNum; + private volatile KafkaTopicManager topicManager; + private volatile RequestStats statsLogger; + private volatile TransactionCoordinator tc; + private volatile String clientHost; + private volatile FetchRequest fetchRequest; + private volatile RequestHeader header; private volatile CompletableFuture resultFuture; - private AtomicBoolean hasComplete; - private AtomicLong bytesReadable; - private DelayedOperationPurgatory fetchPurgatory; - private String namespacePrefix; + private volatile DelayedOperationPurgatory fetchPurgatory; + private volatile String namespacePrefix; // recycler and get for this object public static MessageFetchContext get(KafkaRequestHandler requestHandler, @@ -107,8 +107,6 @@ public static MessageFetchContext get(KafkaRequestHandler requestHandler, String namespacePrefix) { MessageFetchContext context = RECYCLER.get(); context.namespacePrefix = namespacePrefix; - context.responseData = new ConcurrentHashMap<>(); - context.decodeResults = new ConcurrentLinkedQueue<>(); context.requestHandler = requestHandler; context.maxReadEntriesNum = requestHandler.getMaxReadEntriesNum(); context.topicManager = requestHandler.getTopicManager(); @@ -118,8 +116,6 @@ public static MessageFetchContext get(KafkaRequestHandler requestHandler, context.fetchRequest = (FetchRequest) kafkaHeaderAndRequest.getRequest(); context.header = kafkaHeaderAndRequest.getHeader(); context.resultFuture = resultFuture; - context.hasComplete = new AtomicBoolean(false); - context.bytesReadable = new AtomicLong(0); context.fetchPurgatory = fetchPurgatory; return context; } @@ -130,8 +126,6 @@ public static MessageFetchContext getForTest(FetchRequest fetchRequest, CompletableFuture resultFuture) { MessageFetchContext context = RECYCLER.get(); context.namespacePrefix = namespacePrefix; - context.responseData = new ConcurrentHashMap<>(); - context.decodeResults = new ConcurrentLinkedQueue<>(); context.requestHandler = null; context.maxReadEntriesNum = 0; context.topicManager = null; @@ -141,7 +135,6 @@ public static MessageFetchContext getForTest(FetchRequest fetchRequest, context.fetchRequest = fetchRequest; context.header = null; context.resultFuture = resultFuture; - context.hasComplete = new AtomicBoolean(false); return context; } @@ -151,8 +144,10 @@ private MessageFetchContext(Handle recyclerHandle) { private void recycle() { - responseData = null; - decodeResults = null; + responseData.clear(); + decodeResults.clear(); + hasComplete.set(false); + bytesReadable.set(0L); requestHandler = null; maxReadEntriesNum = 0; topicManager = null; @@ -162,8 +157,6 @@ private void recycle() { fetchRequest = null; header = null; resultFuture = null; - hasComplete = null; - bytesReadable = null; fetchPurgatory = null; namespacePrefix = null; recyclerHandle.recycle(this);