New to Kestra?
Use blueprints to kickstart your first workflows.
Validate NATS JetStream events with DuckDB, filter payload anomalies, and publish quarantined records with Slack notifications.
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.
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.
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.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.check_messages (io.kestra.plugin.core.flow.If): Evaluates outputs.consume_events.messagesCount > 0. If false, logs a notice and skips processing.convert_ion_to_json (io.kestra.plugin.serdes.json.IonToJson): Converts internal Ion storage records into newline-delimited JSON for DuckDB query consumption.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).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.publish_valid_events (io.kestra.plugin.core.flow.If): When valid_events.size > 0, publishes validated records to NATS subject events.processed.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.log_summary (io.kestra.plugin.core.log.Log): Logs total consumed, valid, and quarantined event counts.errors: Contains a SlackIncomingWebhook task to alert SREs if any workflow task fails unexpectedly.nats-js).events.>.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
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)nats_url input (default: nats://nats-js:4222).seed_demo_events: true for initial validation run.events.processed, quarantined events in events.quarantine, and Slack channel notifications.io.kestra.plugin.jdbc.postgresql.Query) after publish_valid_events to persist clean records directly into a data warehouse.events.quarantine, fixes payload schemas, and re-publishes to events.sensor.events.processed or events.quarantine, published messages are silently dropped.events.sensor to prevent infinite feedback loops with events.processed.MALFORMED_JSON without failing the flow execution.