Producer metrics¶
Warning
The producer metrics API is experimental and may change without a deprecation period.
AIOKafkaProducer can emit lightweight producer batch lifecycle metrics via
the metrics_collector constructor argument. The collector is a synchronous
callback object derived from ProducerMetricsCollector; aiokafka does not
aggregate, sample, export, or depend on a metrics backend.
Implementations should keep callbacks fast and non-blocking. If you need to do
asynchronous work, push the event into an internal queue and process it from a
separate task. Exceptions raised by collectors are logged and ignored.
Every callback has a no-op default, so implementations only need to override
the callbacks they use. Defining an unknown method whose name starts with
on_ emits a warning to help catch callback name typos.
from aiokafka import AIOKafkaProducer, ProducerMetricsCollector
class CompletionMetrics(ProducerMetricsCollector):
def on_batch_completed(
self,
*,
topic,
partition,
send_to_completion_seconds,
record_count,
acknowledged,
batch_age_seconds,
):
print(
topic,
partition,
record_count,
acknowledged,
batch_age_seconds,
)
producer = AIOKafkaProducer(
bootstrap_servers="localhost:9092",
metrics_collector=CompletionMetrics(),
)
Callbacks¶
on_batch_dispatched(*, topic, partition, queue_time_seconds, batch_size_bytes, record_count, attempt)Called every time a batch leaves the accumulator and is handed to the sender. The first attempt is
1. Retried batches emit this callback again with an incremented attempt.queue_time_secondscovers only the time spent in the accumulator for this attempt. Retry backoff and earlier attempts are excluded.batch_size_bytesandrecord_countare therefore observed per attempt. Count successfully completed records in a terminal callback instead of this callback if retries must not be double-counted.on_batch_completed(*, topic, partition, send_to_completion_seconds, record_count, acknowledged, batch_age_seconds)Called once when the producer successfully completes a batch.
send_to_completion_secondsmeasures time from the final dispatch until producer-side completion.acknowledgedisTruewhen a broker acknowledgment was received. Withacks=0it isFalseand completion means that the producer finished sending the request without waiting for a broker response.batch_age_secondsmeasures time from batch creation to completion and includes retries.on_batch_failed(*, topic, partition, exception, record_count, batch_age_seconds)Called once when a batch ultimately fails.
batch_age_secondsmeasures time from batch creation to failure and includes retries.on_buffer_wait(*, topic, partition, wait_seconds)Called once after a
send()orsend_batch()operation that waited for accumulator space.wait_secondsis accumulated across all wait loops in that operation. It is also reported when the operation ultimately times out or otherwise fails after waiting.
Granularity¶
The dispatch, completion, and failure callbacks are batch-level events.
on_buffer_wait is an operation-level event for a single send() or
send_batch() call.
batch_age_seconds is an upper bound for per-message latency because records
can be appended after their batch was created. The API does not track append
timestamps for individual records and therefore cannot provide exact
per-message end-to-end latency.
Empty batches emit no lifecycle metrics.
All durations are reported in seconds. topic and partition are always
passed; collector implementations decide whether to use them as labels.
Prometheus example¶
from aiokafka import ProducerMetricsCollector
from prometheus_client import Counter, Histogram
queue_time = Histogram(
"aiokafka_producer_batch_queue_time_seconds",
"Time batches spent queued per dispatch attempt.",
["topic", "partition"],
)
send_to_completion = Histogram(
"aiokafka_producer_send_to_completion_seconds",
"Time from final batch dispatch to producer-side completion.",
["topic", "partition", "acknowledged"],
)
batch_age = Histogram(
"aiokafka_producer_batch_age_seconds",
"Time from batch creation to terminal outcome.",
["topic", "partition", "outcome"],
)
batch_size = Histogram(
"aiokafka_producer_batch_size_bytes",
"Producer batch size in bytes per dispatch attempt.",
["topic", "partition"],
)
batch_retries = Counter(
"aiokafka_producer_batch_retries_total",
"Producer batch retry attempts.",
["topic", "partition"],
)
records_completed = Counter(
"aiokafka_producer_records_completed_total",
"Records successfully completed by the producer.",
["topic", "partition", "acknowledged"],
)
batch_failures = Counter(
"aiokafka_producer_batch_failures_total",
"Producer batches that ultimately failed.",
["topic", "partition", "exception"],
)
buffer_wait = Counter(
"aiokafka_producer_buffer_wait_seconds_total",
"Total time spent waiting for producer accumulator space.",
["topic", "partition"],
)
class PrometheusProducerMetrics(ProducerMetricsCollector):
def on_batch_dispatched(
self,
*,
topic,
partition,
queue_time_seconds,
batch_size_bytes,
record_count,
attempt,
):
queue_time.labels(topic, partition).observe(queue_time_seconds)
batch_size.labels(topic, partition).observe(batch_size_bytes)
if attempt > 1:
batch_retries.labels(topic, partition).inc()
def on_batch_completed(
self,
*,
topic,
partition,
send_to_completion_seconds,
record_count,
acknowledged,
batch_age_seconds,
):
acknowledged_label = str(acknowledged).lower()
outcome = "acknowledged" if acknowledged else "unacknowledged"
send_to_completion.labels(
topic, partition, acknowledged_label
).observe(send_to_completion_seconds)
records_completed.labels(
topic, partition, acknowledged_label
).inc(record_count)
batch_age.labels(topic, partition, outcome).observe(batch_age_seconds)
def on_batch_failed(
self,
*,
topic,
partition,
exception,
record_count,
batch_age_seconds,
):
batch_failures.labels(
topic, partition, type(exception).__name__
).inc()
batch_age.labels(topic, partition, "failed").observe(batch_age_seconds)
def on_buffer_wait(self, *, topic, partition, wait_seconds):
buffer_wait.labels(topic, partition).inc(wait_seconds)
API reference¶
- class aiokafka.metrics.ProducerMetricsCollector[source]
Receive experimental producer batch lifecycle events.
All methods are synchronous and called from the producer hot path. Implementations should be fast and non-blocking. If asynchronous work is required, defer it via an internal queue.
Exceptions raised by collectors are logged and ignored.
Each callback has a no-op default implementation. Subclasses only need to override the callbacks they use. Unknown methods starting with
on_emit a warning when the subclass is defined to help detect callback name typos.All time-valued arguments are in seconds.
Warning
This API is experimental and may change without a deprecation period.
- on_batch_dispatched(*, topic: str, partition: int, queue_time_seconds: float, batch_size_bytes: int, record_count: int, attempt: int) None[source]
Called each time a batch is handed to the sender.
queue_time_secondscovers only the queue time for this attempt.attemptstarts at 1. Batch size and record count are reported per attempt and can therefore be observed more than once for retried batches.
- on_batch_completed(*, topic: str, partition: int, send_to_completion_seconds: float, record_count: int, acknowledged: bool, batch_age_seconds: float) None[source]
Called when the producer successfully completes a batch.
acknowledgedis false when the producer is configured withacks=0.batch_age_secondsis an upper bound for the latency of records in the batch and includes retries.