Skip to content
This repository was archived by the owner on Jan 24, 2024. It is now read-only.
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -83,21 +83,21 @@ protected MessageFetchContext newObject(Handle<MessageFetchContext> handle) {
};

private final Handle<MessageFetchContext> recyclerHandle;
private Map<TopicPartition, PartitionData<MemoryRecords>> responseData;
private ConcurrentLinkedQueue<DecodeResult> 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<TopicPartition, PartitionData<MemoryRecords>> responseData = new ConcurrentHashMap<>();
private final ConcurrentLinkedQueue<DecodeResult> 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<AbstractResponse> resultFuture;
private AtomicBoolean hasComplete;
private AtomicLong bytesReadable;
private DelayedOperationPurgatory<DelayedOperation> fetchPurgatory;
private String namespacePrefix;
private volatile DelayedOperationPurgatory<DelayedOperation> fetchPurgatory;
private volatile String namespacePrefix;

// recycler and get for this object
public static MessageFetchContext get(KafkaRequestHandler requestHandler,
Expand All @@ -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();
Expand All @@ -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;
}
Expand All @@ -130,8 +126,6 @@ public static MessageFetchContext getForTest(FetchRequest fetchRequest,
CompletableFuture<AbstractResponse> 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;
Expand All @@ -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;
}

Expand All @@ -151,8 +144,10 @@ private MessageFetchContext(Handle<MessageFetchContext> recyclerHandle) {


private void recycle() {
responseData = null;
decodeResults = null;
responseData.clear();
decodeResults.clear();
hasComplete.set(false);
bytesReadable.set(0L);
requestHandler = null;
maxReadEntriesNum = 0;
topicManager = null;
Expand All @@ -162,8 +157,6 @@ private void recycle() {
fetchRequest = null;
header = null;
resultFuture = null;
hasComplete = null;
bytesReadable = null;
fetchPurgatory = null;
namespacePrefix = null;
recyclerHandle.recycle(this);
Expand Down