[ 
https://issues.apache.org/jira/browse/KAFKA-21098?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Avishek Das updated KAFKA-21098:
--------------------------------
    Description: 
h2. Summary

On a tiered-storage broker, the per-topic {{RemoteDeleteLagSegments}} (and its 
byte twin {{{}RemoteDeleteLagBytes{}}}) gauge can get latched at a non-zero 
value that never drains back to 0, even though remote deletion is healthy and 
there is no orphaned segment in remote storage. The phantom value appears when 
a partition leader moves to another broker (any leadership change — planned 
reassignment, preferred-leader election, or a rebalance) while an expiration 
pass for that partition is still in flight on the old leader. It is purely a 
metrics artifact: there is no real backlog and no data loss. Frequent 
leadership movement makes the broker/cluster aggregate step up over time as 
each affected partition contributes a stuck residual.
h2. Affects

Any tiered-storage-enabled topic in 3.9.0 through trunk (reproduced on 4.2.1). 
The likelihood scales with leadership-change frequency: clusters that reassign 
or rebalance leadership often (e.g. automated balancers, rolling restarts, or 
under-provisioned brokers where leadership shifts a lot) accumulate more stuck 
residuals.
h2. Root cause

The copy path was hardened for this exact race (KAFKA-16948: reset lag metrics 
on becoming follower, and guard {{recordLagStats}} with 
{{{}!isCancelled(){}}}), but the delete path was never given the equivalent 
treatment. Two gaps combine:

(1) *No cancellation guard on the emit.* 
{{RLMExpirationTask.updateRemoteDeleteLagWith(...)}} writes 
{{{}RemoteDeleteLagSegments{}}}/{{{}RemoteDeleteLagBytes{}}} unconditionally. 
There is no {{if (!isCancelled())}} check — unlike the copy path's 
{{{}recordLagStats{}}}, which already skips emitting for a cancelled task.

(2) *The removal and the in-flight emit race.* On a leader→follower transition 
{{onLeadershipChange}} runs {{doHandleFollowerPartition}} (which cancels the 
{{{}RLMExpirationTask{}}}) and then {{removeRemoteTopicPartitionMetrics}} 
(which removes the delete-lag gauges). But an expiration-pool thread that is 
already mid-{{{}cleanupExpiredRemoteLogSegments{}}} for that partition can call 
{{updateRemoteDeleteLagWith(n, bytes)}} *after* the gauge has been removed, 
re-registering it at a non-zero value.

Because the same pass then reaches {{{}deleteRemoteLogSegment(..., ignored -> 
!isCancelled()){}}}, which returns {{false}} for the now-cancelled task, the 
segment is *not* deleted by this old leader and the counter is never 
decremented back down. The gauge is therefore pinned at the re-registered value 
until the same broker re-leads that partition and completes a fresh expiration 
pass, the partition is stopped, or the broker restarts.

Note the underlying segment is *not* orphaned: the new leader independently 
re-evaluates retention and deletes it. So the only symptom is a stuck metric / 
false-positive alert.

h2. Worked example

Config: a tiered-storage topic ({{remote.storage.enable=true}}) with short 
retention so expiration runs continuously, and a setup where leadership for one 
of its partitions moves off broker A to broker B mid-pass.

*Phase A — steady state on broker A (leader):* each expiration cycle sets 
{{RemoteDeleteLagSegments}} to the batch of retention-breached segments at pass 
start, then decrements to 0 as each segment is deleted. Healthy sawtooth, floor 
= 0.

*Phase B — leadership moves off broker A mid-pass:*
|| Step || Thread || Action || Delete-lag gauge on A ||
| 1 | expiration pool | pass starts, finds 1 breached segment, 
{{updateRemoteDeleteLagWith(1, sz)}} | 1 |
| 2 | leadership handler | {{onLeadershipChange}} → cancel 
{{RLMExpirationTask}} for the partition | 1 |
| 3 | leadership handler | {{removeRemoteTopicPartitionMetrics}} removes the 
gauge | (removed) |
| 4 | expiration pool | still in the same pass, {{updateRemoteDeleteLagWith(1, 
sz)}} runs again | 1 (re-registered) |
| 5 | expiration pool | {{deleteRemoteLogSegment(..., ignored -> 
!isCancelled())}} returns false → no delete, no decrement | 1 (stuck) |

The gauge on broker A now stays at 1 forever. Under frequent leadership churn 
the broker/cluster aggregate steps up monotonically as each quiet partition 
contributes its stuck 1.

h2. Steps to reproduce

# Enable tiered storage on a topic with short retention so the expiration task 
deletes segments regularly.
# Drive enough produce traffic that segments roll and become retention-breached.
# Repeatedly move leadership of the partition to another broker while 
expiration is active (reassignment / preferred-leader election / rebalance).
# Observe 
{{kafka.server:type=BrokerTopicMetrics,name=RemoteDeleteLagSegments,topic=<topic>}}
 on the former leader: it stays pinned at a non-zero value with no 
{{RemoteDeleteErrors}} and no {{RemoteStorageException}}, and does not drain to 
0.

h2. Expected vs actual

* *Expected:* {{RemoteDeleteLagSegments}}/{{RemoteDeleteLagBytes}} return to 0 
on a broker once it is no longer the leader for the partition (the new leader 
owns deletion), mirroring the copy-lag behavior fixed in KAFKA-16948.
* *Actual:* the metric stays latched at a non-zero value (typically 1) on the 
former leader indefinitely, producing false-positive "delete lag stuck" alerts.


h2. Proposed fix

Mirror the copy path. Guard the emit with the task's cancellation state so a 
cancelled expiration task cannot re-register the gauge after 
{{removeRemoteTopicPartitionMetrics}} has cleared it:

{code:java}
// VisibleForTesting
void updateRemoteDeleteLagWith(int segmentsLeftToDelete, long 
sizeOfDeletableSegmentsBytes) {
    if (!isCancelled()) {
        String topic = topicIdPartition.topic();
        int partition = topicIdPartition.partition();
        brokerTopicStats.recordRemoteDeleteLagSegments(topic, partition, 
segmentsLeftToDelete);
        brokerTopicStats.recordRemoteDeleteLagBytes(topic, partition, 
sizeOfDeletableSegmentsBytes);
    }
}
{code}

This closes the race in the same way KAFKA-16948 did for {{recordLagStats}} on 
the copy path. Fix proposed in PR https://github.com/apache/kafka/pull/23471.

h2. Related issues

* KAFKA-16948 — reset lag metrics on becoming follower (guarded the *copy* 
path's {{recordLagStats}}); this bug is the un-fixed *delete* path equivalent.
* KAFKA-15147 — measure pending and outstanding remote segment operations 
(umbrella that introduced the RemoteCopyLag/RemoteDeleteLag metrics).
* KAFKA-20977 — a separate phantom-lag off-by-one on the copy path 
(single-record segments); different root cause, same "gauge stuck non-zero" 
symptom class.



  was:
h2. Summary

On a tiered-storage broker, the per-topic {{RemoteDeleteLagSegments}} (and its 
byte twin {{{}RemoteDeleteLagBytes{}}}) gauge can get latched at a non-zero 
value that never drains back to 0, even though remote deletion is healthy and 
there is no orphaned segment in remote storage. The phantom value appears when 
a partition leader moves to another broker (any leadership change — planned 
reassignment, preferred-leader election, or a rebalance) while an expiration 
pass for that partition is still in flight on the old leader. It is purely a 
metrics artifact: there is no real backlog and no data loss. Frequent 
leadership movement makes the broker/cluster aggregate step up over time as 
each affected partition contributes a stuck residual.
h2. Affects

Any tiered-storage-enabled topic in 3.9.0 through trunk (reproduced on 4.2.1). 
The likelihood scales with leadership-change frequency: clusters that reassign 
or rebalance leadership often (e.g. automated balancers, rolling restarts, or 
under-provisioned brokers where leadership shifts a lot) accumulate more stuck 
residuals.
h2. Root cause

The copy path was hardened for this exact race (KAFKA-16948: reset lag metrics 
on becoming follower, and guard {{recordLagStats}} with 
{{{}!isCancelled(){}}}), but the delete path was never given the equivalent 
treatment. Two gaps combine:

(1) *No cancellation guard on the emit.* 
{{RLMExpirationTask.updateRemoteDeleteLagWith(...)}} writes 
{{{}RemoteDeleteLagSegments{}}}/{{{}RemoteDeleteLagBytes{}}} unconditionally. 
There is no {{if (!isCancelled())}} check — unlike the copy path's 
{{{}recordLagStats{}}}, which already skips emitting for a cancelled task.

(2) *The removal and the in-flight emit race.* On a leader→follower transition 
{{onLeadershipChange}} runs {{doHandleFollowerPartition}} (which cancels the 
{{{}RLMExpirationTask{}}}) and then {{removeRemoteTopicPartitionMetrics}} 
(which removes the delete-lag gauges). But an expiration-pool thread that is 
already mid-{{{}cleanupExpiredRemoteLogSegments{}}} for that partition can call 
{{updateRemoteDeleteLagWith(n, bytes)}} *after* the gauge has been removed, 
re-registering it at a non-zero value.

Because the same pass then reaches {{{}deleteRemoteLogSegment(..., ignored -> 
!isCancelled()){}}}, which returns {{false}} for the now-cancelled task, the 
segment is *not* deleted by this old leader and the counter is never 
decremented back down. The gauge is therefore pinned at the re-registered value 
until the same broker re-leads that partition and completes a fresh expiration 
pass, the partition is stopped, or the broker restarts.

Note the underlying segment is *not* orphaned: the new leader independently 
re-evaluates retention and deletes it. So the only symptom is a stuck metric / 
false-positive alert.

h2. Worked example

Config: a tiered-storage topic ({{remote.storage.enable=true}}) with short 
retention so expiration runs continuously, and a setup where leadership for one 
of its partitions moves off broker A to broker B mid-pass.

*Phase A — steady state on broker A (leader):* each expiration cycle sets 
{{RemoteDeleteLagSegments}} to the batch of retention-breached segments at pass 
start, then decrements to 0 as each segment is deleted. Healthy sawtooth, floor 
= 0.

*Phase B — leadership moves off broker A mid-pass:*
|| Step || Thread || Action || Delete-lag gauge on A ||
| 1 | expiration pool | pass starts, finds 1 breached segment, 
{{updateRemoteDeleteLagWith(1, sz)}} | 1 |
| 2 | leadership handler | {{onLeadershipChange}} → cancel 
{{RLMExpirationTask}} for the partition | 1 |
| 3 | leadership handler | {{removeRemoteTopicPartitionMetrics}} removes the 
gauge | (removed) |
| 4 | expiration pool | still in the same pass, {{updateRemoteDeleteLagWith(1, 
sz)}} runs again | 1 (re-registered) |
| 5 | expiration pool | {{deleteRemoteLogSegment(..., ignored -> 
!isCancelled())}} returns false → no delete, no decrement | 1 (stuck) |

The gauge on broker A now stays at 1 forever. Under frequent leadership churn 
the broker/cluster aggregate steps up monotonically as each quiet partition 
contributes its stuck 1.

h2. Steps to reproduce

# Enable tiered storage on a topic with short retention so the expiration task 
deletes segments regularly.
# Drive enough produce traffic that segments roll and become retention-breached.
# Repeatedly move leadership of the partition to another broker while 
expiration is active (reassignment / preferred-leader election / rebalance).
# Observe 
{{kafka.server:type=BrokerTopicMetrics,name=RemoteDeleteLagSegments,topic=<topic>}}
 on the former leader: it stays pinned at a non-zero value with no 
{{RemoteDeleteErrors}} and no {{RemoteStorageException}}, and does not drain to 
0.

h2. Expected vs actual

* *Expected:* {{RemoteDeleteLagSegments}}/{{RemoteDeleteLagBytes}} return to 0 
on a broker once it is no longer the leader for the partition (the new leader 
owns deletion), mirroring the copy-lag behavior fixed in KAFKA-16948.
* *Actual:* the metric stays latched at a non-zero value (typically 1) on the 
former leader indefinitely, producing false-positive "delete lag stuck" alerts.


h2. Proposed fix

Mirror the copy path. Guard the emit with the task's cancellation state so a 
cancelled expiration task cannot re-register the gauge after 
{{removeRemoteTopicPartitionMetrics}} has cleared it:

{code:java}
// VisibleForTesting
void updateRemoteDeleteLagWith(int segmentsLeftToDelete, long 
sizeOfDeletableSegmentsBytes) {
    if (!isCancelled()) {
        String topic = topicIdPartition.topic();
        int partition = topicIdPartition.partition();
        brokerTopicStats.recordRemoteDeleteLagSegments(topic, partition, 
segmentsLeftToDelete);
        brokerTopicStats.recordRemoteDeleteLagBytes(topic, partition, 
sizeOfDeletableSegmentsBytes);
    }
}
{code}

This closes the race in the same way KAFKA-16948 did for {{recordLagStats}} on 
the copy path. Fix proposed in PR #<TBD>.

h2. Related issues

* KAFKA-16948 — reset lag metrics on becoming follower (guarded the *copy* 
path's {{recordLagStats}}); this bug is the un-fixed *delete* path equivalent.
* KAFKA-15147 — measure pending and outstanding remote segment operations 
(umbrella that introduced the RemoteCopyLag/RemoteDeleteLag metrics).
* KAFKA-20977 — a separate phantom-lag off-by-one on the copy path 
(single-record segments); different root cause, same "gauge stuck non-zero" 
symptom class.




> RemoteDeleteLagSegments/RemoteDeleteLagBytes report a phantom lag that never 
> drains after a leader→follower transition
> ----------------------------------------------------------------------------------------------------------------------
>
>                 Key: KAFKA-21098
>                 URL: https://issues.apache.org/jira/browse/KAFKA-21098
>             Project: Kafka
>          Issue Type: Bug
>          Components: Tiered-Storage
>    Affects Versions: 3.9.2, 4.0.2, 4.1.2, 4.2.1, 4.3.1
>            Reporter: Avishek Das
>            Assignee: Avishek Das
>            Priority: Minor
>              Labels: tiered-storage
>
> h2. Summary
> On a tiered-storage broker, the per-topic {{RemoteDeleteLagSegments}} (and 
> its byte twin {{{}RemoteDeleteLagBytes{}}}) gauge can get latched at a 
> non-zero value that never drains back to 0, even though remote deletion is 
> healthy and there is no orphaned segment in remote storage. The phantom value 
> appears when a partition leader moves to another broker (any leadership 
> change — planned reassignment, preferred-leader election, or a rebalance) 
> while an expiration pass for that partition is still in flight on the old 
> leader. It is purely a metrics artifact: there is no real backlog and no data 
> loss. Frequent leadership movement makes the broker/cluster aggregate step up 
> over time as each affected partition contributes a stuck residual.
> h2. Affects
> Any tiered-storage-enabled topic in 3.9.0 through trunk (reproduced on 
> 4.2.1). The likelihood scales with leadership-change frequency: clusters that 
> reassign or rebalance leadership often (e.g. automated balancers, rolling 
> restarts, or under-provisioned brokers where leadership shifts a lot) 
> accumulate more stuck residuals.
> h2. Root cause
> The copy path was hardened for this exact race (KAFKA-16948: reset lag 
> metrics on becoming follower, and guard {{recordLagStats}} with 
> {{{}!isCancelled(){}}}), but the delete path was never given the equivalent 
> treatment. Two gaps combine:
> (1) *No cancellation guard on the emit.* 
> {{RLMExpirationTask.updateRemoteDeleteLagWith(...)}} writes 
> {{{}RemoteDeleteLagSegments{}}}/{{{}RemoteDeleteLagBytes{}}} unconditionally. 
> There is no {{if (!isCancelled())}} check — unlike the copy path's 
> {{{}recordLagStats{}}}, which already skips emitting for a cancelled task.
> (2) *The removal and the in-flight emit race.* On a leader→follower 
> transition {{onLeadershipChange}} runs {{doHandleFollowerPartition}} (which 
> cancels the {{{}RLMExpirationTask{}}}) and then 
> {{removeRemoteTopicPartitionMetrics}} (which removes the delete-lag gauges). 
> But an expiration-pool thread that is already 
> mid-{{{}cleanupExpiredRemoteLogSegments{}}} for that partition can call 
> {{updateRemoteDeleteLagWith(n, bytes)}} *after* the gauge has been removed, 
> re-registering it at a non-zero value.
> Because the same pass then reaches {{{}deleteRemoteLogSegment(..., ignored -> 
> !isCancelled()){}}}, which returns {{false}} for the now-cancelled task, the 
> segment is *not* deleted by this old leader and the counter is never 
> decremented back down. The gauge is therefore pinned at the re-registered 
> value until the same broker re-leads that partition and completes a fresh 
> expiration pass, the partition is stopped, or the broker restarts.
> Note the underlying segment is *not* orphaned: the new leader independently 
> re-evaluates retention and deletes it. So the only symptom is a stuck metric 
> / false-positive alert.
> h2. Worked example
> Config: a tiered-storage topic ({{remote.storage.enable=true}}) with short 
> retention so expiration runs continuously, and a setup where leadership for 
> one of its partitions moves off broker A to broker B mid-pass.
> *Phase A — steady state on broker A (leader):* each expiration cycle sets 
> {{RemoteDeleteLagSegments}} to the batch of retention-breached segments at 
> pass start, then decrements to 0 as each segment is deleted. Healthy 
> sawtooth, floor = 0.
> *Phase B — leadership moves off broker A mid-pass:*
> || Step || Thread || Action || Delete-lag gauge on A ||
> | 1 | expiration pool | pass starts, finds 1 breached segment, 
> {{updateRemoteDeleteLagWith(1, sz)}} | 1 |
> | 2 | leadership handler | {{onLeadershipChange}} → cancel 
> {{RLMExpirationTask}} for the partition | 1 |
> | 3 | leadership handler | {{removeRemoteTopicPartitionMetrics}} removes the 
> gauge | (removed) |
> | 4 | expiration pool | still in the same pass, 
> {{updateRemoteDeleteLagWith(1, sz)}} runs again | 1 (re-registered) |
> | 5 | expiration pool | {{deleteRemoteLogSegment(..., ignored -> 
> !isCancelled())}} returns false → no delete, no decrement | 1 (stuck) |
> The gauge on broker A now stays at 1 forever. Under frequent leadership churn 
> the broker/cluster aggregate steps up monotonically as each quiet partition 
> contributes its stuck 1.
> h2. Steps to reproduce
> # Enable tiered storage on a topic with short retention so the expiration 
> task deletes segments regularly.
> # Drive enough produce traffic that segments roll and become 
> retention-breached.
> # Repeatedly move leadership of the partition to another broker while 
> expiration is active (reassignment / preferred-leader election / rebalance).
> # Observe 
> {{kafka.server:type=BrokerTopicMetrics,name=RemoteDeleteLagSegments,topic=<topic>}}
>  on the former leader: it stays pinned at a non-zero value with no 
> {{RemoteDeleteErrors}} and no {{RemoteStorageException}}, and does not drain 
> to 0.
> h2. Expected vs actual
> * *Expected:* {{RemoteDeleteLagSegments}}/{{RemoteDeleteLagBytes}} return to 
> 0 on a broker once it is no longer the leader for the partition (the new 
> leader owns deletion), mirroring the copy-lag behavior fixed in KAFKA-16948.
> * *Actual:* the metric stays latched at a non-zero value (typically 1) on the 
> former leader indefinitely, producing false-positive "delete lag stuck" 
> alerts.
> h2. Proposed fix
> Mirror the copy path. Guard the emit with the task's cancellation state so a 
> cancelled expiration task cannot re-register the gauge after 
> {{removeRemoteTopicPartitionMetrics}} has cleared it:
> {code:java}
> // VisibleForTesting
> void updateRemoteDeleteLagWith(int segmentsLeftToDelete, long 
> sizeOfDeletableSegmentsBytes) {
>     if (!isCancelled()) {
>         String topic = topicIdPartition.topic();
>         int partition = topicIdPartition.partition();
>         brokerTopicStats.recordRemoteDeleteLagSegments(topic, partition, 
> segmentsLeftToDelete);
>         brokerTopicStats.recordRemoteDeleteLagBytes(topic, partition, 
> sizeOfDeletableSegmentsBytes);
>     }
> }
> {code}
> This closes the race in the same way KAFKA-16948 did for {{recordLagStats}} 
> on the copy path. Fix proposed in PR 
> https://github.com/apache/kafka/pull/23471.
> h2. Related issues
> * KAFKA-16948 — reset lag metrics on becoming follower (guarded the *copy* 
> path's {{recordLagStats}}); this bug is the un-fixed *delete* path equivalent.
> * KAFKA-15147 — measure pending and outstanding remote segment operations 
> (umbrella that introduced the RemoteCopyLag/RemoteDeleteLag metrics).
> * KAFKA-20977 — a separate phantom-lag off-by-one on the copy path 
> (single-record segments); different root cause, same "gauge stuck non-zero" 
> symptom class.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to