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]