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
Choose the outcome
Consumer lag monitoring
- 2
Define the contract
topic, partition, consumer_group, and status
- 3
Log the final result
Emit aggregate lag snapshots separately from per-message completion events at high volume.
- 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.
npm installation
npm install telemetry-sh- 1Create one reusable server-side client. Set its timeout and retry limit.
- 2Log an event when the operation succeeds, fails, retries, or times out.
- 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 consumer lag and processing analytics event
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 consumer lag and processing analytics verification query
SELECT *
FROM kafka_consumer_snapshot
ORDER BY timestamp_utc DESC
LIMIT 20;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
Event schemas for this workflow
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.
Measure event ingestion freshness
Which production event sources are stale or delayed right now?
Open recipeMeasure late-arriving events
Which event producers deliver data late enough to distort analysis?
Open recipeTrack event schema-version adoption
Which producers still emit old versions of a critical event?
Open recipeDetect missing service heartbeats
Which expected telemetry sources have stopped sending heartbeats?
Open recipeDetect error-rate spikes with a rolling baseline
Which hourly error-rate buckets are far above their recent baseline?
Open recipeMeasure Telemetry volume by event name
Which event contracts create the most ingestion volume?
Open recipeBrowse by implementation family
Compare related integration patterns
Templates to pair with this integration
More integrations
Redis and node-redis telemetry
Measure Redis command outcomes, latency, cache behavior, reconnects, and bounded error categories alongside node-redis.
Open guiden8n workflow telemetry
Use the n8n HTTP Request node to send approved fields for workflow completion, failure, retries, item counts, and downstream delivery results to Telemetry.
Open guideAWS Lambda structured event monitoring
Track Lambda invocations, cold starts, duration, retries, and business outcomes with compact structured events.
Open guide