Skip to content
Integration guide

Kafka consumer lag and processing analytics

Monitor Kafka consumer outcomes, partition lag, processing latency, retries, and poison-message handling with structured events.

Reviewed by the Telemetry product team on . We checked which events to send, which data to exclude, and how to add the code. Who reviews this page

Useful for
  • Consumer lag monitoring
  • Partition imbalance
  • Message processing reliability
Test the integration

Send and verify events with Kafka consumer lag and processing analytics

Use Kafka consumer lag and processing analytics where your app knows the final result. Collect only the fields you need, then verify a test event before building charts.

  1. 1

    Choose the outcome

    Consumer lag monitoring

  2. 2

    Define the contract

    topic, partition, consumer_group, and status

  3. 3

    Log the final result

    Emit aggregate lag snapshots separately from per-message completion events at high volume.

  4. 4

    Check the stored event

    Exercise a known fixture, then inspect kafka_consumer_snapshot for one correctly typed terminal row.

Before you start

Before you start

  • A server-side TELEMETRY_API_KEY
  • Stable topic and consumer-group names
  • Offset and high-watermark metadata from the Kafka client

Delivery setup

Install and initialize server-side

Import telemetry-sh in server-only code and initialize it once with process.env.TELEMETRY_API_KEY. Keep ingestion credentials out of browser bundles, client-visible environment variables, source control, logs, and exception messages.

kafka-install

npm installation

bash
npm install telemetry-sh
  1. 1Create one reusable server-side client. Set its timeout and retry limit.
  2. 2Log an event when the operation succeeds, fails, retries, or times out.
  3. 3Send test events with known results and inspect the stored rows before enabling an alert.

Snippet

Start with one structured event

Add this shape where the workflow completes, fails, or retries. Then build the dashboard from real fields.

kafka

Kafka consumer lag and processing analytics event

javascript
await telemetry.log("kafka_consumer_snapshot", {
  topic,
  partition,
  consumer_group: groupId,
  offset: Number(currentOffset),
  high_watermark: Number(highWatermark),
  lag: Number(highWatermark - currentOffset),
  status: "healthy",
  consumer_instance: instanceId,
  release: process.env.APP_RELEASE,
});

Event schema

topic, partition, consumer_group, and status

offset, high_watermark, lag, and duration_ms

attempt, error_type, consumer_instance, and release

Check your setup

Checkpoint 1

Emit aggregate lag snapshots separately from per-message completion events at high volume.

Checkpoint 2

Do not send message values, keys containing customer data, or authentication configuration.

Checkpoint 3

Preserve topic and consumer group while limiting partition-level charts to operational investigation.

Verification

Prove the event arrived

Run this after exercising known success and failure cases. Replace the fallback table name if your final event contract differs from the snippet.

kafka-verification

Kafka consumer lag and processing analytics verification query

sql
SELECT *
FROM kafka_consumer_snapshot
ORDER BY timestamp_utc DESC
LIMIT 20;
Confirm one terminal row per logical outcome, with the expected status, identifiers, units, and UTC time.
Inspect the inferred schema and verify that retries do not change field types or generate a new logical event ID.
Search the stored fields for credentials, raw payloads, prompts, private content, and unbounded error messages.
Exercise a provider timeout, ingestion rejection, and process shutdown before treating the dashboard as complete.

Implementation references

Review the event contract, data-safety guidance, and upstream primary documentation before enabling a new production path.

Where to log

Keep the outcome event small and recoverable

This pattern provides

  • Record the outcome as an event you can query with SQL.
  • Stable fields for dashboards, alerts, and cross-event correlation.
  • Test events for checking success, failure, retries, and timeouts.

This pattern does not provide

  • An OTLP exporter, automatic collection pipeline, or replacement for detailed traces and diagnostic logs.
  • Exactly-once delivery merely because the payload contains an event ID.
  • Permission to collect raw provider payloads, user content, credentials, or regulated data.

Example event schemas

Check what each event records, when to send it, and which field types it needs. Review the example payload and privacy checklist before using it in production.

Use these queries in Telemetry

Learn about Structured events

Send events with consistent names and field types. Choose which context to include before sending it.

Related SQL recipes

More SQL recipes

Run the query using this workflow's event fields and check the example result. Save the result to a dashboard or set up an alert.

Browse all recipes

Browse by implementation family

Compare related integration patterns

Templates to pair with this integration

More integrations