This is an automated email from the ASF dual-hosted git repository.
chia7712 pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git
The following commit(s) were added to refs/heads/trunk by this push:
new 17f3f8aeae8 MINOR: Cleanups following KAFKA-20637 (#22584)
17f3f8aeae8 is described below
commit 17f3f8aeae8c7a93fce731b016ffb055c2bee4bc
Author: Mickael Maison <[email protected]>
AuthorDate: Wed Jul 29 04:56:57 2026 +0200
MINOR: Cleanups following KAFKA-20637 (#22584)
Following up comment on https://github.com/apache/kafka/pull/22408
Reviewers: Chia-Ping Tsai <[email protected]>
---
.../main/scala/kafka/server/ControllerApis.scala | 2 +-
.../scala/kafka/server/RequestHandlerHelper.scala | 4 +--
.../scala/kafka/tools/TestRaftRequestHandler.scala | 6 ++--
.../unit/kafka/network/SocketServerTest.scala | 4 +--
.../java/org/apache/kafka/network/Request.java | 4 +++
.../org/apache/kafka/server/ForwardingManager.java | 32 ++++------------------
.../apache/kafka/server/ForwardingManagerImpl.java | 25 ++++++++---------
.../org/apache/kafka/server/EnvelopeUtilsTest.java | 2 +-
8 files changed, 29 insertions(+), 50 deletions(-)
diff --git a/core/src/main/scala/kafka/server/ControllerApis.scala
b/core/src/main/scala/kafka/server/ControllerApis.scala
index 9f56906f822..cc42cba56b8 100644
--- a/core/src/main/scala/kafka/server/ControllerApis.scala
+++ b/core/src/main/scala/kafka/server/ControllerApis.scala
@@ -686,7 +686,7 @@ class ControllerApis(
request: Request,
buildResponse: ApiMessage => AbstractResponse
): CompletableFuture[Unit] = {
- val requestBody = request.body(classOf[AbstractRequest])
+ val requestBody = request.body
val future = raftManager
.handleRequest(request.context, request.header, requestBody.data,
time.milliseconds())
.toCompletableFuture
diff --git a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
index 29ee72e025b..9c78451d674 100644
--- a/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
+++ b/core/src/main/scala/kafka/server/RequestHandlerHelper.scala
@@ -20,7 +20,7 @@ package kafka.server
import kafka.network.RequestChannel
import kafka.server.QuotaFactory.QuotaManagers
import org.apache.kafka.common.errors.ClusterAuthorizationException
-import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse}
+import org.apache.kafka.common.requests.AbstractResponse
import org.apache.kafka.common.utils.Time
import org.apache.kafka.network.Request
import org.apache.kafka.server.quota.{ClientQuotaManager,
ControllerMutationQuota, ThrottleCallback}
@@ -56,7 +56,7 @@ class RequestHandlerHelper(
error: Throwable,
throttleMs: Int
): Unit = {
- val requestBody = request.body(classOf[AbstractRequest])
+ val requestBody = request.body
val response = requestBody.getErrorResponse(throttleMs, error)
if (response == null)
requestChannel.closeConnection(request, requestBody.errorCounts(error))
diff --git a/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
b/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
index 28badb26c3f..d9646e8dff9 100644
--- a/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
+++ b/core/src/main/scala/kafka/tools/TestRaftRequestHandler.scala
@@ -23,7 +23,7 @@ import kafka.utils.Logging
import org.apache.kafka.common.internals.FatalExitError
import org.apache.kafka.common.message.{BeginQuorumEpochResponseData,
EndQuorumEpochResponseData, FetchResponseData, FetchSnapshotResponseData,
VoteResponseData}
import org.apache.kafka.common.protocol.{ApiKeys, ApiMessage}
-import org.apache.kafka.common.requests.{AbstractRequest, AbstractResponse,
BeginQuorumEpochResponse, EndQuorumEpochResponse, FetchResponse,
FetchSnapshotResponse, VoteResponse}
+import org.apache.kafka.common.requests.{AbstractResponse,
BeginQuorumEpochResponse, EndQuorumEpochResponse, FetchResponse,
FetchSnapshotResponse, VoteResponse}
import org.apache.kafka.common.utils.Time
import org.apache.kafka.network.Request
import org.apache.kafka.raft.RaftManager
@@ -57,7 +57,7 @@ class TestRaftRequestHandler(
case e: Throwable =>
error(s"Unexpected error handling request ${request.requestDesc(true)}
" +
s"with context ${request.context}", e)
- val errorResponse =
request.body(classOf[AbstractRequest]).getErrorResponse(e)
+ val errorResponse = request.body.getErrorResponse(e)
requestChannel.sendResponse(request, errorResponse)
} finally {
// The local completion time may be set while processing the request.
Only record it if it's unset.
@@ -94,7 +94,7 @@ class TestRaftRequestHandler(
request: Request,
buildResponse: ApiMessage => AbstractResponse
): Unit = {
- val requestBody = request.body(classOf[AbstractRequest])
+ val requestBody = request.body
val future = raftManager.handleRequest(
request.context,
diff --git a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
index 51ea50628ef..bc4cfe4a513 100644
--- a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
+++ b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala
@@ -161,7 +161,7 @@ class SocketServerTest {
}
def processRequest(channel: RequestChannel, request: Request): Unit = {
- val byteBuffer =
request.body(classOf[AbstractRequest]).serializeWithHeader(request.header)
+ val byteBuffer = request.body.serializeWithHeader(request.header)
val send = new NetworkSend(request.context.connectionId,
ByteBufferSend.sizePrefixed(byteBuffer))
val headerLog = RequestConvertToJson.requestHeaderNode(request.header)
channel.sendResponse(new SendResponse(request, send,
Optional.of(headerLog)))
@@ -641,7 +641,7 @@ class SocketServerTest {
// Mimic a primitive request handler that fetches the request from
RequestChannel and place a response with a
// throttled channel.
val request = receiveRequest(server.dataPlaneRequestChannel)
- val byteBuffer =
request.body(classOf[AbstractRequest]).serializeWithHeader(request.header)
+ val byteBuffer = request.body.serializeWithHeader(request.header)
val send = new NetworkSend(request.context.connectionId,
ByteBufferSend.sizePrefixed(byteBuffer))
val channelThrottlingCallback = new ThrottleCallback {
diff --git a/server/src/main/java/org/apache/kafka/network/Request.java
b/server/src/main/java/org/apache/kafka/network/Request.java
index c9144889496..e9258e4e3a1 100644
--- a/server/src/main/java/org/apache/kafka/network/Request.java
+++ b/server/src/main/java/org/apache/kafka/network/Request.java
@@ -282,6 +282,10 @@ public final class Request implements BaseRequest {
+ ", but found " + bodyAndSize.request.getClass().getName());
}
+ public AbstractRequest body() {
+ return bodyAndSize.request;
+ }
+
public AbstractRequest loggableRequest() {
if (bodyAndSize.request instanceof AlterConfigsRequest alterConfigs) {
var newData = alterConfigs.data().duplicate();
diff --git
a/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
b/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
index 79e26566929..708e137c3a6 100644
--- a/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
+++ b/server/src/main/java/org/apache/kafka/server/ForwardingManager.java
@@ -19,20 +19,13 @@ package org.apache.kafka.server;
import org.apache.kafka.clients.NodeApiVersions;
import org.apache.kafka.common.requests.AbstractRequest;
import org.apache.kafka.common.requests.AbstractResponse;
-import org.apache.kafka.common.requests.RequestContext;
import org.apache.kafka.network.Request;
import java.nio.ByteBuffer;
import java.util.Optional;
import java.util.function.Consumer;
-import java.util.function.Supplier;
-public interface ForwardingManager {
-
- /**
- * Close the forwarding manager
- */
- void close();
+public interface ForwardingManager extends AutoCloseable {
/**
* Forward given request to the active controller.
@@ -49,12 +42,7 @@ public interface ForwardingManager {
) {
ByteBuffer buffer = originalRequest.buffer().duplicate();
buffer.flip();
- forwardRequest(originalRequest.context(),
- buffer,
- originalRequest.startTimeNanos(),
- originalRequest.body(AbstractRequest.class),
- originalRequest::toString,
- responseCallback);
+ forwardRequest(originalRequest, buffer, originalRequest.body(),
responseCallback);
}
/**
@@ -72,36 +60,26 @@ public interface ForwardingManager {
AbstractRequest newRequestBody,
Consumer<Optional<AbstractResponse>> responseCallback) {
ByteBuffer buffer =
newRequestBody.serializeWithHeader(originalRequest.header());
- forwardRequest(originalRequest.context(),
- buffer,
- originalRequest.startTimeNanos(),
- newRequestBody,
- originalRequest::toString,
- responseCallback);
+ forwardRequest(originalRequest, buffer, newRequestBody,
responseCallback);
}
/**
* Forward given request to the active controller.
*
- * @param requestContext The request context of the original envelope
request.
+ * @param originalRequest The request to forward.
* @param requestBufferCopy The request buffer we want to send. This
should not be the original
* byte buffer from the envelope request, since
we will be mutating
* the position and limit fields. It should be
a copy.
- * @param requestCreationNs The request creation timestamp in
nanoseconds.
* @param requestBody The AbstractRequest we are sending.
- * @param requestToString A callback which can be invoked to produce a
human-readable description
- * of the request.
* @param responseCallback A callback which takes in an
`Optional<AbstractResponse>`.
* We will call this function with
Optional.of(x) after the controller responds with x.
* Or, if the controller doesn't support the
request version, we will complete
* the callback with Optional.empty().
*/
void forwardRequest(
- RequestContext requestContext,
+ Request originalRequest,
ByteBuffer requestBufferCopy,
- long requestCreationNs,
AbstractRequest requestBody,
- Supplier<String> requestToString,
Consumer<Optional<AbstractResponse>> responseCallback);
/**
diff --git
a/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
b/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
index 0344bd67d74..63cc4b54572 100644
--- a/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
+++ b/server/src/main/java/org/apache/kafka/server/ForwardingManagerImpl.java
@@ -25,8 +25,8 @@ import org.apache.kafka.common.requests.AbstractRequest;
import org.apache.kafka.common.requests.AbstractResponse;
import org.apache.kafka.common.requests.EnvelopeRequest;
import org.apache.kafka.common.requests.EnvelopeResponse;
-import org.apache.kafka.common.requests.RequestContext;
import org.apache.kafka.common.requests.RequestHeader;
+import org.apache.kafka.network.Request;
import org.apache.kafka.server.common.ControllerRequestCompletionHandler;
import org.apache.kafka.server.common.NodeToControllerChannelManager;
import org.apache.kafka.server.metrics.ForwardingManagerMetrics;
@@ -38,9 +38,8 @@ import java.nio.ByteBuffer;
import java.util.Optional;
import java.util.concurrent.TimeUnit;
import java.util.function.Consumer;
-import java.util.function.Supplier;
-public class ForwardingManagerImpl implements ForwardingManager, AutoCloseable
{
+public class ForwardingManagerImpl implements ForwardingManager {
private static final Logger LOG =
LoggerFactory.getLogger(ForwardingManagerImpl.class);
@@ -58,14 +57,12 @@ public class ForwardingManagerImpl implements
ForwardingManager, AutoCloseable {
@Override
public void forwardRequest(
- RequestContext requestContext,
+ Request originalRequest,
ByteBuffer requestBufferCopy,
- long requestCreationNs,
AbstractRequest requestBody,
- Supplier<String> requestToString,
Consumer<Optional<AbstractResponse>> responseCallback) {
- EnvelopeRequest.Builder envelopeRequest =
ForwardingManagerUtil.buildEnvelopeRequest(requestContext, requestBufferCopy);
- long requestCreationTimeMs =
TimeUnit.NANOSECONDS.toMillis(requestCreationNs);
+ EnvelopeRequest.Builder envelopeRequest =
ForwardingManagerUtil.buildEnvelopeRequest(originalRequest.context(),
requestBufferCopy);
+ long requestCreationTimeMs =
TimeUnit.NANOSECONDS.toMillis(originalRequest.startTimeNanos());
class ForwardingResponseHandler implements
ControllerRequestCompletionHandler {
@@ -77,11 +74,11 @@ public class ForwardingManagerImpl implements
ForwardingManager, AutoCloseable {
if (clientResponse.versionMismatch() != null) {
LOG.debug("Returning `UNKNOWN_SERVER_ERROR` in response to
{} due to unexpected version error",
- requestToString.get(),
clientResponse.versionMismatch());
+ originalRequest, clientResponse.versionMismatch());
responseCallback.accept(Optional.of(requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception())));
} else if (clientResponse.authenticationException() != null) {
LOG.debug("Returning `UNKNOWN_SERVER_ERROR` in response to
{} due to authentication error",
- requestToString.get(),
clientResponse.authenticationException());
+ originalRequest,
clientResponse.authenticationException());
responseCallback.accept(Optional.of(requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception())));
} else {
EnvelopeResponse envelopeResponse = (EnvelopeResponse)
clientResponse.responseBody();
@@ -102,10 +99,10 @@ public class ForwardingManagerImpl implements
ForwardingManager, AutoCloseable {
// return `UNKNOWN_SERVER_ERROR` so that the user
knows that there is a problem
// on the broker.
LOG.debug("Forwarded request {} failed with an
error in the envelope response {}",
- requestToString.get(), envelopeError);
+ originalRequest, envelopeError);
response =
requestBody.getErrorResponse(Errors.UNKNOWN_SERVER_ERROR.exception());
} else {
- response =
parseResponse(envelopeResponse.responseData(), requestBody,
requestContext.header);
+ response =
parseResponse(envelopeResponse.responseData(), requestBody,
originalRequest.context().header);
}
responseCallback.accept(Optional.of(response));
}
@@ -114,7 +111,7 @@ public class ForwardingManagerImpl implements
ForwardingManager, AutoCloseable {
@Override
public void onTimeout() {
- LOG.debug("Forwarding of the request {} failed due to timeout
exception", requestToString.get());
+ LOG.debug("Forwarding of the request {} failed due to timeout
exception", originalRequest);
forwardingManagerMetrics.decrementQueueLength();
forwardingManagerMetrics.queueTimeMsHist().record(channelManager.getTimeoutMs());
AbstractResponse response = requestBody.getErrorResponse(new
TimeoutException());
@@ -136,7 +133,7 @@ public class ForwardingManagerImpl implements
ForwardingManager, AutoCloseable {
return channelManager.controllerApiVersions();
}
- private AbstractResponse parseResponse(ByteBuffer buffer, AbstractRequest
request, RequestHeader header) {
+ private static AbstractResponse parseResponse(ByteBuffer buffer,
AbstractRequest request, RequestHeader header) {
try {
return AbstractResponse.parseResponse(buffer, header);
} catch (Exception e) {
diff --git
a/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
b/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
index e365fab9575..af7e6351bc2 100644
--- a/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
+++ b/server/src/test/java/org/apache/kafka/server/EnvelopeUtilsTest.java
@@ -86,7 +86,7 @@ public class EnvelopeUtilsTest {
assertTrue(forwardedRequest.isForwarded());
assertSame(envelope, forwardedRequest.envelope().orElseThrow());
assertEquals(envelope.requestDequeueTimeNanos(),
forwardedRequest.requestDequeueTimeNanos());
- assertInstanceOf(CreateTopicsRequest.class,
forwardedRequest.body(AbstractRequest.class));
+ assertInstanceOf(CreateTopicsRequest.class, forwardedRequest.body());
}
@Test