Eric Chang created KAFKA-21031:
----------------------------------
Summary: Another busy loop happens while a heartbeat is in flight
and the heartbeat timer is expired
Key: KAFKA-21031
URL: https://issues.apache.org/jira/browse/KAFKA-21031
Project: Kafka
Issue Type: Bug
Components: clients, consumer
Reporter: Eric Chang
Assignee: Eric Chang
Fix For: 4.4.0, 4.5.0
The async consumer (\{{group.protocol=consumer}}), the share consumer and the
streams consumer busy-spin on both the network thread and the application
thread while a group heartbeat request (\{{ConsumerGroupHeartbeat}},
{{ShareGroupHeartbeat}} or {{{}StreamsGroupHeartbeat{}}}) is in flight and the
heartbeat timer has expired. Both threads spin until the response or a failure
arrives.
h4. How it happens
{{HeartbeatRequestState}} resets the heartbeat timer when a heartbeat request
is generated and marked in flight (it may still be queued for an unreachable
host). While that request is in flight the manager cannot send another one, so
{{AbstractHeartbeatRequestManager.poll()}} and {{maximumTimeToWait()}} fall
back to {{{}timeToNextHeartbeatMs(){}}}. If the timer is expired at that point,
it returns {{{}remainingBackoffMs(){}}}, which is measured from the last
response and, with default settings, is already 0. {{poll()}} then returns
{{PollResult(0)}} and {{ConsumerNetworkThread}} calls {{client.poll(0)}} in a
tight loop; {{maximumTimeToWait()}} returns 0 and the application thread does
not block in {{poll()}} either, waking the network thread on every iteration.
The timer is expired while a request is in flight in two cases:
# The first heartbeat of a consumer instance. {{heartbeatIntervalMs}} is
initialised to 0 and is only known from the first heartbeat response, so
{{resetTimer()}} leaves the timer expired immediately.
# # Any later heartbeat whose response takes longer than the heartbeat
interval (for example a coordinator in a long GC pause). The interval is known,
but the timer expires with the request still in flight.
The guards added by KAFKA-20253 and KAFKA-20970 / KAFKA-21010 check
{{{}coordinator().isEmpty(){}}}; they do not apply because the coordinator is
already known.
h4. Impact
* Consumer restarts during a coordinator move (rolling broker upgrade): the
first heartbeat goes to a coordinator that accepts the connection but does not
answer, or is not reachable yet. Both threads spin until {{request.timeout.ms}}
(30 s), or until the connection attempt fails if the host is not reachable. The
failure then marks the coordinator unknown and the normal jittered, exponential
retry and coordinator rediscovery apply; the next heartbeat spins again until
it is answered.
* * Coordinator overloaded or in a long GC pause: any consumer whose heartbeat
response is slower than the broker-provided heartbeat interval (5 s by default)
spins for the rest of the pause.
* * Coordinator not reachable from the client (wrong
{{{}advertised.listeners{}}}, firewall): the consumer never recovers and
repeats the cycle above. Request and connection timeouts are logged at
INFO/WARN, but nothing identifies the busy loop itself.
* * Every normal start: a few milliseconds of spin per instance.
h4. Reproduction
{{HeartbeatRequestState}} with {{{}heartbeatIntervalMs = 0{}}},
{{{}onSendAttempt(){}}}, {{{}resetTimer(){}}}, then {{timeToNextHeartbeatMs()}}
returns 0 while {{canSendRequest()}} is false. Same with interval 5000 after
{{onSuccessfulAttempt()}} and advancing the clock past the interval. Driving
{{AbstractHeartbeatRequestManager.poll()}} the same way shows
{{PollResult.timeUntilNextPollMs == 0}} with no request sent, and
{{{}maximumTimeToWait() == 0{}}}.
h4. History
The async-consumer first-heartbeat case is present since 3.7.0: KAFKA-15278 (PR
14364) introduced the zero initial interval and KAFKA-15890 (PR 14835)
introduced {{{}maximumTimeToWait(){}}}. KAFKA-16528 (PR 15698), first released
in 3.8.0, moved the timer reset from response time to send time; since then a
later heartbeat enters this state only when its response takes longer than the
interval. The share and streams consumers inherited the same state machine in
3.9.0 and 4.1.0. Nothing in 4.4 or 4.5 makes it more likely; it was found by
auditing the wait-time paths after KAFKA-20970. It is the same class of bug as
KAFKA-20253, which was backported to 4.2.2 and 4.3.2. Trunk is already
4.5.0-SNAPSHOT and the 4.4 branch has been cut; since 4.4 is still in the RC
phase, the fix is small and has no behaviour change in the steady state,
including it in the next 4.4 RC is proposed, with cherry-picks to 4.3 and 4.2
as for KAFKA-20253.
h4. Proposed fix
{{HeartbeatRequestState.timeToNextHeartbeatMs()}} returns the effective initial
retry backoff (\{{min(retry.backoff.ms, retry.backoff.max.ms)}}, 100 ms by
default) when the timer is expired and a request is in flight, the same value
KAFKA-21010 (PR 23348, pending) proposes for the coordinator-unknown guard. The
two fixes are complementary: a request timeout moves the client from this
known-coordinator, in-flight state to the unknown-coordinator state handled by
KAFKA-21010. Waiting forever is not safe because {{NetworkClient}} only notices
a timed-out request after its poll returns and the network thread caps that
poll at 5 s. Both backoff configs accept 0, so the value is floored at 1 ms. PR
to follow.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)