Schedule icon
If icon
Produce icon
Consume icon
IonToJson icon
Query icon
SlackIncomingWebhook icon
Log icon

Validate NATS JetStream events with DuckDB and quarantine failures

Validate NATS JetStream events with DuckDB, filter payload anomalies, and publish quarantined records with Slack notifications.

Categories
DataInfrastructure

Diagram unavailable

We could not build the topology for this blueprint. The flow itself is valid, use the YAML on the left to run it.

Validate incoming NATS JetStream sensor events using DuckDB in-memory SQL processing and route non-compliant records to a quarantine topic with Slack alerts.

Who it's for

  • Data Engineers and Event-Driven Architecture Operators who require automated quality gates, schema enforcement, and payload validation for real-time messaging streams.
  • Infrastructure teams managing NATS JetStream messaging clusters who need inline stream cleansing before writing records to downstream databases.
  • System Reliability Engineers who need real-time Slack notifications when telemetry streams contain payload anomalies or malformed JSON messages.

Why orchestrate this with Kestra

NATS JetStream ingests streaming data at high velocity, but it has no native schema validation or SQL query engine. Hardcoding schema validation inside downstream microservices leads to duplicated code, fragile exception handling, and silent data corruption. Kestra solves this by orchestrating the entire lifecycle: consuming batches from JetStream via a durable consumer, running zero-setup DuckDB SQL quality checks in memory, routing clean records to events.processed, writing invalid payloads to events.quarantine, and alerting engineers over Slack.

How it works

  1. seed_events (io.kestra.plugin.core.flow.If): When seed_demo_events input is true, publishes 7 sample sensor events (2 valid, 5 invalid/malformed) to events.sensor using io.kestra.plugin.nats.core.Produce.
  2. consume_events (io.kestra.plugin.nats.core.Consume): Pulls a batch of messages from events.sensor using durable consumer kestra-event-validator and writes raw records to Kestra internal storage in Ion format.
  3. check_messages (io.kestra.plugin.core.flow.If): Evaluates outputs.consume_events.messagesCount > 0. If false, logs a notice and skips processing.
  4. convert_ion_to_json (io.kestra.plugin.serdes.json.IonToJson): Converts internal Ion storage records into newline-delimited JSON for DuckDB query consumption.
  5. valid_events (io.kestra.plugin.jdbc.duckdb.Query): Executes DuckDB SQL classification query to extract rows where reason IS NULL (valid event_id, sensor_id, temperature within range, unit celsius).
  6. quarantined_events (io.kestra.plugin.jdbc.duckdb.Query): Executes DuckDB SQL classification query to extract rows where reason IS NOT NULL (MALFORMED_JSON, MISSING_EVENT_ID, MISSING_SENSOR_ID, NULL_TEMPERATURE, TEMPERATURE_OUT_OF_RANGE, INVALID_UNIT), outputting JSON objects containing the error reason and original payload.
  7. publish_valid_events (io.kestra.plugin.core.flow.If): When valid_events.size > 0, publishes validated records to NATS subject events.processed.
  8. handle_quarantine_events (io.kestra.plugin.core.flow.If): When quarantined_events.size > 0, publishes quarantined records to NATS subject events.quarantine and sends a summary notification via io.kestra.plugin.slack.notifications.SlackIncomingWebhook.
  9. log_summary (io.kestra.plugin.core.log.Log): Logs total consumed, valid, and quarantined event counts.
  10. errors: Contains a SlackIncomingWebhook task to alert SREs if any workflow task fails unexpectedly.

Prerequisites

  • A running NATS JetStream server (e.g., container named nats-js).
  • A NATS Stream configured to capture subjects matching events.>.
  • Connectivity between the Kestra worker container and the NATS server network.
  • A Slack Incoming Webhook URL configured in Kestra secrets.

Setup Commands

Start a NATS JetStream container on your Kestra network:

docker run -d --name nats-js --network <kestra_network> -p 4222:4222 -p 8222:8222 nats:latest -js

Create JetStream stream for events.> using nats-box CLI container:

docker run --rm --network <kestra_network> natsio/nats-box nats -s nats://nats-js:4222 stream add SENSOR_EVENTS --subjects "events.>" --storage file --defaults

Secrets

  • SLACK_WEBHOOK_URL: Slack Incoming Webhook URL. (Note for open-source users: Store secrets in environment variables prefixed with SECRET_ encoded in base64, e.g. SECRET_SLACK_WEBHOOK_URL)

Quick start

  1. Configure nats_url input (default: nats://nats-js:4222).
  2. Keep seed_demo_events: true for initial validation run.
  3. Execute the flow and inspect valid records in events.processed, quarantined events in events.quarantine, and Slack channel notifications.

How to extend

  • Add custom DuckDB classification branches for additional sensor parameters like pressure, humidity, or battery percentage.
  • Add a Postgres or Snowflake sink task (io.kestra.plugin.jdbc.postgresql.Query) after publish_valid_events to persist clean records directly into a data warehouse.
  • Replace Slack webhook notifications with PagerDuty alerts or Microsoft Teams webhooks using Kestra plugin tasks.
  • Configure automated dead-letter queue replay by building a secondary flow that reads from events.quarantine, fixes payload schemas, and re-publishes to events.sensor.

Links

Pitfalls

  • At-Least-Once Acknowledgment: Messages are acknowledged upon consume fetch. If a downstream task fails, the consumed batch is not redelivered automatically.
  • Core NATS Publishing: The Produce task publishes using core NATS. If no JetStream stream captures events.processed or events.quarantine, published messages are silently dropped.
  • Stream Requirement: The target stream must exist in NATS before executing the Consume task.
  • Exact Subject Filter: The consumer filters specifically on events.sensor to prevent infinite feedback loops with events.processed.
  • Malformed Payload Handling: Non-JSON payloads land in quarantine with MALFORMED_JSON without failing the flow execution.

What you get

  • A single, fully automated event quality pipeline that consumes, validates, and routes NATS JetStream messages without custom microservice code.
  • Zero-copy SQL validation powered by DuckDB in-memory engine, performing schema enforcement directly on streaming payloads.
  • Complementary partitioning of event batches into valid payloads and quarantined anomaly records with detailed root cause tags.
  • Immediate Slack alerting for quarantine spikes and execution failure tracking across your message bus.
See How

New to Kestra?

Use blueprints to kickstart your first workflows.