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:

Reply via email to