This is an automated email from the ASF dual-hosted git repository.
lucasbru 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 ca0c89d83ae MINOR: Add logging for topology description push in
streams group protocol (#22670)
ca0c89d83ae is described below
commit ca0c89d83ae20a3f684170c965bb2ae4b750b533
Author: Lucas Brutschy <[email protected]>
AuthorDate: Thu Jun 25 14:43:54 2026 +0200
MINOR: Add logging for topology description push in streams group protocol
(#22670)
Add INFO log lines for the topology description push flow in the streams
group protocol: when the broker requests a push via heartbeat response,
when the push request is sent, and when it completes successfully.
Reviewers: TengYao Chi <[email protected]>, Alieh Saeedi
<[email protected]>
---
.../consumer/internals/StreamsGroupHeartbeatRequestManager.java | 1 +
.../internals/StreamsGroupTopologyDescriptionRequestManager.java | 3 +++
2 files changed, 4 insertions(+)
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
index fce71fe95d5..6ae3e048c7c 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupHeartbeatRequestManager.java
@@ -572,6 +572,7 @@ public class StreamsGroupHeartbeatRequestManager implements
RequestManager {
streamsRebalanceData.setAcceptableRecoveryLag(data.acceptableRecoveryLag());
if (data.topologyDescriptionRequired() &&
streamsRebalanceData.wireTopologyDescription() != null) {
+ logger.info("Broker requested topology description push");
streamsRebalanceData.setTopologyPushRequired(true);
}
diff --git
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupTopologyDescriptionRequestManager.java
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupTopologyDescriptionRequestManager.java
index 30979905693..c2da1c7b2f4 100644
---
a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupTopologyDescriptionRequestManager.java
+++
b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/StreamsGroupTopologyDescriptionRequestManager.java
@@ -72,6 +72,8 @@ public class StreamsGroupTopologyDescriptionRequestManager
implements RequestMan
.setTopologyEpoch(streamsRebalanceData.topologyEpoch())
.setTopologyDescription(streamsRebalanceData.wireTopologyDescription());
+ logger.info("Sending topology description for group {}",
data.groupId());
+
final NetworkClientDelegate.UnsentRequest unsent = new
NetworkClientDelegate.UnsentRequest(
new StreamsGroupTopologyDescriptionUpdateRequest.Builder(data),
coordinatorRequestManager.coordinator()
@@ -140,6 +142,7 @@ public class StreamsGroupTopologyDescriptionRequestManager
implements RequestMan
case NONE:
pushRequestState.onSuccessfulAttempt(responseTimeMs);
streamsRebalanceData.setTopologyPushRequired(false);
+ logger.info("Topology description pushed successfully");
break;
case NOT_COORDINATOR: