DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
Status
Current state: Under Discussion
Discussion thread: https://lists.apache.org/thread/z4o5093zqxmc18nzb74vyhvb1hb399v4
JIRA:
KAFKA-18088
-
Getting issue details...
STATUS
Please keep the discussion on the mailing list rather than commenting on the wiki (wiki discussions get unwieldy fast).
Motivation
The Kafka consumer allows users to `pause()` and `resume()` individual partitions. When a partition is paused, subsequent calls to `poll()` will not return any records from that partition until it is resumed. This mechanism is used for application-level flow control and backpressure management.
Higher-level frameworks like Kafka Streams make extensive use of pause/resume internally:
- Backpressure: When an internal record buffer for a partition exceeds its capacity, Kafka Streams pauses that partition to prevent further fetching (`StreamTask.addRecords()`). It resumes the partition once the buffer drains below the threshold (`StreamTask.resumePollingForPartitionsWithAvailableSpace()`).
- Rebalance lifecycle: After a rebalance, Kafka Streams pauses all partitions that are not yet owned by fully initialized tasks (`TaskManager`). Partitions are only resumed after state store restoration completes.
- Changelog restore management: The `StoreChangelogReader` pauses and resumes changelog partitions on the restore consumer depending on whether the corresponding tasks still exist.
Currently, there are no consumer metrics that expose pause/resume state. Operators and developers have no way to answer basic questions through monitoring:
- How many partitions are currently paused? A number of paused partitions may indicate backpressure issues, slow state restoration, or application-level processing bottlenecks.
- Which specific partitions are paused? Knowing which partitions are paused helps pinpoint which topics or partitions are experiencing issues.
- How long has a partition been paused? A partition that has been paused for an extended period may indicate a stuck consumer, an unrecoverable error in processing, or a state store restoration that is taking too long.
This KIP proposes adding new consumer metrics that track paused partitions and their pause duration.
Public Interfaces
New metrics
consumer-metrics
The metric will have the following tags:
- client-id
| Metric Name | Type | Recording Level | Description |
| paused-partitions-count | Gauge (Integer) | INFO | The current number of partitions that have been paused by the user via `Consumer.pause()`. |
consumer-fetch-manager-metrics
Each metric will have the following tags:
- client-id
- topic
- partition
| Metric Name | Type | Recording Level | Description |
| paused-partitions | Gauge (Integer) | DEBUG | Whether this partition is currently paused. Returns `1` if paused, `0` if not paused. |
| paused-partitions-duration-seconds | Gauge (Long) | DEBUG | The time in seconds since this partition was paused. Returns `-1` if the partition is not currently paused. |
Proposed Changes
Track Pause Timestamp
Currently, each partition's internal state tracks whether it is paused with a simple boolean flag. To support the `paused-partitions-duration-seconds` metric, this KIP adds a timestamp that records when the partition was paused. The timestamp is set when pause() is called, cleared (reset to -1) when resume() is called, and also reset to -1 on partition reassignment regardless of prior pause status. This ensures that a partition re-assigned to the same consumer after a rebalance starts with a clean pause state.
consumer-metrics
A `paused-partitions-count` gauge is registered in the consumer's top-level metrics. When queried, this gauge computes the current count of paused partitions by reading from the consumer's subscription state. Since the consumer already maintains the set of paused partitions internally, no additional data structures are needed for this metric.
consumer-fetch-manager-metrics
The `paused-partitions` and `paused-partitions-duration-seconds` gauges are registered per partition in the fetch manager metrics, alongside the existing per-partition lag and lead metrics. They follow the same lifecycle:
- On partition assignment: The metrics are registered for each newly assigned partition.
- On partition revocation: The metrics are removed for each revoked partition, preventing stale metrics from accumulating. Both metrics are implemented as on-demand gauges that read directly from the partition's pause state when queried, rather than being updated on every pause/resume call.
Compatibility, Deprecation, and Migration Plan
This change only adds new metrics. No existing metrics or APIs are deprecated.
Test Plan
- Verify that the pause timestamp is set when `pause()` is called, cleared on `resume()`, and reset on partition reassignment.
- Verify that the `paused-partitions-count` metric is registered and returns the correct count as partitions are paused and resumed.
- Verify that `paused-partitions` and `paused-partitions-duration-seconds` metrics are registered per partition, return correct values, and are removed when partitions are revoked.
Rejected Alternatives
N/A