kevin-wu24 commented on code in PR #22111:
URL: https://github.com/apache/kafka/pull/22111#discussion_r3482216725


##########
raft/src/test/java/org/apache/kafka/raft/KafkaRaftClientFetchTest.java:
##########
@@ -765,4 +766,122 @@ void testUpdatedHighWatermarkCompleted() throws Exception 
{
             assertEquals(localLogEndOffset, partitionResponse.highWatermark());
         }
     }
+
+    @Test
+    void testObserverFetchesBetweenLeaderAndBootstrapServers() throws 
Exception {

Review Comment:
   Here is an updated trace from my most recent local changes:
   ```
   [2026-06-26 09:50:15,417] INFO Starting request manager with bootstrap 
servers: [localhost:10139 (id: -2 rack: null isFenced: false)] 
(org.apache.kafka.raft.KafkaRaftClient:331)
   [2026-06-26 09:50:15,600] INFO Reading KRaft snapshot and log as part of the 
initialization (org.apache.kafka.raft.KafkaRaftClient:509)
   [2026-06-26 09:50:15,601] INFO Starting voters are 
VoterSet(voters={148=VoterNode(voterKey=ReplicaKey(id=148, 
directoryId=<undefined>), 
listeners=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 supportedKRaftVersion=SupportedVersionRange[min_version:0, max_version:0]), 
149=VoterNode(voterKey=ReplicaKey(id=149, directoryId=<undefined>), 
listeners=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10139}),
 supportedKRaftVersion=SupportedVersionRange[min_version:0, max_version:0])}) 
(org.apache.kafka.raft.KafkaRaftClient:511)
   [2026-06-26 09:50:15,603] INFO Attempting durable transition to 
UnattachedState(epoch=0, leaderId=OptionalInt.empty, votedKey=Optional.empty, 
voters=[148, 149], electionTimeoutMs=18985, highWatermark=Optional.empty) from 
null (org.apache.kafka.raft.QuorumState:732)
   [2026-06-26 09:50:15,605] INFO Completed transition to 
UnattachedState(epoch=0, leaderId=OptionalInt.empty, votedKey=Optional.empty, 
voters=[148, 149], electionTimeoutMs=18985, highWatermark=Optional.empty) from 
null (org.apache.kafka.raft.QuorumState:744)
   [2026-06-26 09:50:15,614] TRACE Sent outbound request: 
OutboundRequest(correlationId=0, 
data=FetchRequestData(clusterId='sSoE9smGSQqjfEuTnlMPsA', replicaId=-1, 
replicaState=ReplicaState(replicaId=147, replicaEpoch=-1), maxWaitMs=0, 
minBytes=0, maxBytes=1048576, isolationLevel=0, sessionId=0, sessionEpoch=-1, 
topics=[FetchTopic(topic='metadata', topicId=AAAAAAAAAAAAAAAAAAAAAQ, 
partitions=[FetchPartition(partition=0, currentLeaderEpoch=0, fetchOffset=0, 
lastFetchedEpoch=0, logStartOffset=-1, partitionMaxBytes=0, 
replicaDirectoryId=XiEwxtuzSGuh5WQsWw8VnQ, highWatermark=-1)])], 
forgottenTopicsData=[], rackId=''), createdTimeMs=1782485415405, 
destination=localhost:10139 (id: -2 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2908)
   [2026-06-26 09:50:15,615] INFO Registered the listener 
org.apache.kafka.raft.RaftClientTestContext$MockListener@220558713 
(org.apache.kafka.raft.KafkaRaftClient:3590)
   [2026-06-26 09:50:15,707] TRACE Received inbound message 
InboundResponse(correlationId=0, data=FetchResponseData(throttleTimeMs=0, 
errorCode=0, sessionId=0, responses=[FetchableTopicResponse(topic='', 
topicId=AAAAAAAAAAAAAAAAAAAAAQ, partitions=[PartitionData(partitionIndex=0, 
errorCode=6, highWatermark=0, lastStableOffset=-1, logStartOffset=-1, 
divergingEpoch=EpochEndOffset(epoch=-1, endOffset=-1), 
currentLeader=LeaderIdAndEpoch(leaderId=148, leaderEpoch=2), 
snapshotId=SnapshotId(endOffset=-1, epoch=-1), abortedTransactions=[], 
preferredReadReplica=-1, records=MemoryRecords(size=0, 
buffer=java.nio.HeapByteBuffer[pos=0 lim=0 cap=37]))])], 
nodeEndpoints=[NodeEndpoint(nodeId=148, host='localhost', port=10138, 
rack=null)]), source=localhost:10139 (id: -2 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2848)
   [2026-06-26 09:50:15,708] INFO Attempting durable transition to 
FollowerState(fetchTimeoutMs=50000, epoch=2, leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) from UnattachedState(epoch=0, 
leaderId=OptionalInt.empty, votedKey=Optional.empty, voters=[148, 149], 
electionTimeoutMs=18985, highWatermark=Optional.empty) 
(org.apache.kafka.raft.QuorumState:732)
   [2026-06-26 09:50:15,709] INFO Completed transition to 
FollowerState(fetchTimeoutMs=50000, epoch=2, leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) from UnattachedState(epoch=0, 
leaderId=OptionalInt.empty, votedKey=Optional.empty, voters=[148, 149], 
electionTimeoutMs=18985, highWatermark=Optional.empty) 
(org.apache.kafka.raft.QuorumState:744)
   [2026-06-26 09:50:15,709] DEBUG Notifying listener 
org.apache.kafka.raft.RaftClientTestContext$MockListener@220558713 of leader 
change LeaderAndEpoch[leaderId=OptionalInt[148], epoch=2] 
(org.apache.kafka.raft.KafkaRaftClient:4121)
   [2026-06-26 09:50:15,812] TRACE Sent outbound request: 
OutboundRequest(correlationId=1, 
data=FetchRequestData(clusterId='sSoE9smGSQqjfEuTnlMPsA', replicaId=-1, 
replicaState=ReplicaState(replicaId=147, replicaEpoch=-1), maxWaitMs=0, 
minBytes=0, maxBytes=1048576, isolationLevel=0, sessionId=0, sessionEpoch=-1, 
topics=[FetchTopic(topic='metadata', topicId=AAAAAAAAAAAAAAAAAAAAAQ, 
partitions=[FetchPartition(partition=0, currentLeaderEpoch=2, fetchOffset=0, 
lastFetchedEpoch=0, logStartOffset=-1, partitionMaxBytes=0, 
replicaDirectoryId=XiEwxtuzSGuh5WQsWw8VnQ, highWatermark=-1)])], 
forgottenTopicsData=[], rackId=''), createdTimeMs=1782485415405, 
destination=localhost:10138 (id: 148 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2908)
   [2026-06-26 09:50:15,813] TRACE Received inbound message 
InboundResponse(correlationId=1, data=FetchResponseData(throttleTimeMs=0, 
errorCode=8, sessionId=0, responses=[], nodeEndpoints=[]), 
source=localhost:10138 (id: 148 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2848)
   [2026-06-26 09:50:15,814] INFO Attempting durable transition to 
UnattachedState(epoch=2, leaderId=OptionalInt[148], votedKey=Optional.empty, 
voters=[148, 149], electionTimeoutMs=9223372036854775807, 
highWatermark=Optional.empty) from FollowerState(fetchTimeoutMs=50000, epoch=2, 
leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) (org.apache.kafka.raft.QuorumState:732)
   [2026-06-26 09:50:15,814] INFO Completed transition to 
UnattachedState(epoch=2, leaderId=OptionalInt[148], votedKey=Optional.empty, 
voters=[148, 149], electionTimeoutMs=9223372036854775807, 
highWatermark=Optional.empty) from FollowerState(fetchTimeoutMs=50000, epoch=2, 
leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) (org.apache.kafka.raft.QuorumState:744)
   [2026-06-26 09:50:15,918] TRACE Sent outbound request: 
OutboundRequest(correlationId=2, 
data=FetchRequestData(clusterId='sSoE9smGSQqjfEuTnlMPsA', replicaId=-1, 
replicaState=ReplicaState(replicaId=147, replicaEpoch=-1), maxWaitMs=0, 
minBytes=0, maxBytes=1048576, isolationLevel=0, sessionId=0, sessionEpoch=-1, 
topics=[FetchTopic(topic='metadata', topicId=AAAAAAAAAAAAAAAAAAAAAQ, 
partitions=[FetchPartition(partition=0, currentLeaderEpoch=2, fetchOffset=0, 
lastFetchedEpoch=0, logStartOffset=-1, partitionMaxBytes=0, 
replicaDirectoryId=XiEwxtuzSGuh5WQsWw8VnQ, highWatermark=-1)])], 
forgottenTopicsData=[], rackId=''), createdTimeMs=1782485465406, 
destination=localhost:10139 (id: -2 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2908)
   [2026-06-26 09:50:15,919] TRACE Received inbound message 
InboundResponse(correlationId=2, data=FetchResponseData(throttleTimeMs=0, 
errorCode=0, sessionId=0, responses=[FetchableTopicResponse(topic='', 
topicId=AAAAAAAAAAAAAAAAAAAAAQ, partitions=[PartitionData(partitionIndex=0, 
errorCode=6, highWatermark=0, lastStableOffset=-1, logStartOffset=-1, 
divergingEpoch=EpochEndOffset(epoch=-1, endOffset=-1), 
currentLeader=LeaderIdAndEpoch(leaderId=148, leaderEpoch=2), 
snapshotId=SnapshotId(endOffset=-1, epoch=-1), abortedTransactions=[], 
preferredReadReplica=-1, records=MemoryRecords(size=0, 
buffer=java.nio.HeapByteBuffer[pos=0 lim=0 cap=37]))])], 
nodeEndpoints=[NodeEndpoint(nodeId=148, host='localhost', port=10138, 
rack=null)]), source=localhost:10139 (id: -2 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2848)
   [2026-06-26 09:50:15,920] INFO Attempting durable transition to 
FollowerState(fetchTimeoutMs=50000, epoch=2, leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) from UnattachedState(epoch=2, 
leaderId=OptionalInt[148], votedKey=Optional.empty, voters=[148, 149], 
electionTimeoutMs=9223372036854775807, highWatermark=Optional.empty) 
(org.apache.kafka.raft.QuorumState:732)
   [2026-06-26 09:50:15,920] INFO Completed transition to 
FollowerState(fetchTimeoutMs=50000, epoch=2, leader=148, 
leaderEndpoints=Endpoints(endpoints={ListenerName(LISTENER)=localhost/<unresolved>:10138}),
 votedKey=Optional.empty, voters=[148, 149], highWatermark=Optional.empty, 
fetchingSnapshot=Optional.empty) from UnattachedState(epoch=2, 
leaderId=OptionalInt[148], votedKey=Optional.empty, voters=[148, 149], 
electionTimeoutMs=9223372036854775807, highWatermark=Optional.empty) 
(org.apache.kafka.raft.QuorumState:744)
   [2026-06-26 09:50:16,024] TRACE Sent outbound request: 
OutboundRequest(correlationId=3, 
data=FetchRequestData(clusterId='sSoE9smGSQqjfEuTnlMPsA', replicaId=-1, 
replicaState=ReplicaState(replicaId=147, replicaEpoch=-1), maxWaitMs=0, 
minBytes=0, maxBytes=1048576, isolationLevel=0, sessionId=0, sessionEpoch=-1, 
topics=[FetchTopic(topic='metadata', topicId=AAAAAAAAAAAAAAAAAAAAAQ, 
partitions=[FetchPartition(partition=0, currentLeaderEpoch=2, fetchOffset=0, 
lastFetchedEpoch=0, logStartOffset=-1, partitionMaxBytes=0, 
replicaDirectoryId=XiEwxtuzSGuh5WQsWw8VnQ, highWatermark=-1)])], 
forgottenTopicsData=[], rackId=''), createdTimeMs=1782485465406, 
destination=localhost:10138 (id: 148 rack: null isFenced: false)) 
(org.apache.kafka.raft.KafkaRaftClient:2908)
   ```
   
   The scenario is: after the local node becomes `Follower`, it is unable to 
successfully fetch from the leader, instead receiving the 
`BROKER_NOT_AVAILABLE` message, for the duration of its fetch timeout. This is 
shown by the local node sending a fetch to the leader, getting a 
`BROKER_NOT_AVAILABLE` response, and only then transitioning to `Unattached`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to