New to Kestra?
Use blueprints to kickstart your first workflows.
Audit Druid sys.segments for published-but-unavailable segments and alert Slack before queries silently return incomplete results.
Druid's own documentation draws a sharp line between two different flags on every
segment: is_published (the metadata store considers it used) and is_available
(some Historical or realtime task is actually serving it right now). A segment can be
published and still unavailable - a Historical crashed mid-reassignment, a load rule
changed before balancing finished, a deep storage blip during startup - and when that
happens, a query against that datasource does not error. It just quietly returns fewer
rows than it should. This blueprint queries sys.segments for exactly that gap,
classifies the cluster's state, and alerts Slack the moment it appears or the cluster
cannot be reached.
druid_segment_audit_schedule (io.kestra.plugin.core.trigger.Schedule) runs every 30 minutes, shipped disabled so you can validate a manual run first.audit_segment_availability (io.kestra.plugin.jdbc.druid.Query) runs SELECT COUNT(*) AS unavailable_count FROM sys.segments WHERE is_published = 1 AND is_available = 0 with fetchType: FETCH_ONE and allowFailure: true, so an unreachable cluster produces no row instead of failing the task outright.classify_availability_risk (io.kestra.plugin.core.debug.Return) classifies the result as UNREACHABLE (no row), AT_RISK (unavailable_count > max_unavailable_segments), or HEALTHY.route_by_availability (io.kestra.plugin.core.flow.Switch) branches on that classification: HEALTHY logs a clean check, AT_RISK pages Slack with the unavailable count, UNREACHABLE pages Slack that the cluster could not be queried, and the defaults: block catches any unexpected classification value.log_audit_status always runs last, printing the classification regardless of which branch fired.errors block alerts Slack separately if the flow itself fails outside these handled branches.HEALTHY/AT_RISK/UNREACHABLE) so a broken audit is never mistaken for a healthy cluster.sys.segments metadata table through the plugin's plain SQL Query task.sys.segments auditing alongside this repo's existing Cassandra system_auth.roles/system_schema.keyspaces sentinels.sys.segments is a live, read-only snapshot; querying it by hand in the Druid console answers one moment in time and notifies nobody. Kestra supplies the schedule, the three-way risk classification as a first-class branch, and an execution history that shows exactly when a published segment first went unavailable and for how long it stayed that way.
Local testing:
docker run -d --name druid-availability-gate -p 8889:8888 -e DRUID_SINGLE_NODE_CONF=micro-quickstart apache/druid:latest
docker network connect YOUR_KESTRA_NETWORK druid-availability-gate (only needed if the Kestra Worker runs in a separate Docker network than this container; set druid_url to use druid-availability-gate as the host from inside that network, or localhost:8889 from the host). Druid's official single-node quickstart image bundles ZooKeeper, the metadata store, and every Druid process in one container for local evaluation; production clusters run these as separate services, matching Druid's own deployment documentation.
SLACK_WEBHOOK_URL: Slack incoming webhook used by alert_at_risk, alert_unreachable, and the errors block.druid_url carries no embedded credentials, matching this plugin's own documented example (druid-to-pandas.yaml), which targets an unauthenticated local Avatica endpoint. For a secured cluster, add username/password properties to audit_segment_availability from {{ secret(...) }} values rather than hardcoding them.SECRET_, base64-encoded, and read back in flows with {{ secret('NAME') }} - for example SECRET_SLACK_WEBHOOK_URL=$(echo -n 'https://hooks.slack.com/...' | base64). This keeps credentials out of the flow YAML but, per Kestra's own documentation, offers no encryption at rest or access control beyond the host environment; use the Enterprise secrets backend for stronger guarantees.druid_url (STRING, default jdbc:avatica:remote:url=http://druid-router:8888/druid/v2/sql/avatica/;transparent_reconnection=true): Avatica JDBC connection string.max_unavailable_segments (INT, default 0): breach threshold; 0 means any published-but-unavailable segment at all breaches.outputs.audit_segment_availability.row.unavailable_count: the raw count for every successful run (absent when the cluster is unreachable).outputs.classify_availability_risk.value: HEALTHY, AT_RISK, or UNREACHABLE for every run.druid_url at an existing cluster.SLACK_WEBHOOK_URL secret.log_healthy fires against a normally-balanced cluster.AT_RISK, stop a Historical process mid-rebalance in a non-production cluster (or lower max_unavailable_segments below the cluster's current transient unavailable count, if any) and re-run; confirm alert_at_risk fires with the expected count.druid_segment_audit_schedule once you trust the check.datasource (SELECT datasource, COUNT(*) ... GROUP BY datasource) to name which datasource is affected in the Slack alert, instead of only a cluster-wide count.druid_url values with io.kestra.plugin.core.flow.Loop (ForEach is deprecated) to sweep a multi-cluster fleet in one run.is_overshadowed = 1 AND is_available = 1 to also catch segments still being served after they should have been superseded, a different kind of staleness risk sys.segments can also surface.sys.servers to report which Historical processes are currently under-serving, for a fuller incident brief alongside the Slack alert.markUsed/markUnused APIs change whether a segment is considered used, not whether it is currently being served, and would be the wrong fix for this specific gap. No safe, verified single-call remediation for "published but unavailable" was found in Druid's own API reference, so this blueprint only audits and alerts.is_available = 0 does not mean the segment is lost. Verified from Druid's own documentation: it only means no Historical or realtime task is serving it right now - often transient during rebalancing, a rolling restart, or cluster scale-down. A single AT_RISK reading is a prompt to check the Coordinator console, not necessarily evidence of data loss.allowFailure: true on audit_segment_availability is what lets UNREACHABLE be classified instead of failing the execution outright, matching this repo's own cassandra-replication-audit-gate.yaml pattern: a connection failure produces an undefined row output rather than aborting the flow, which classify_availability_risk reads as UNREACHABLE.is_published/is_available/is_overshadowed are returned as 0/1 longs, not SQL booleans. Verified from Druid's own SQL metadata table reference; the = 1 / = 0 comparisons in this blueprint's query match that documented representation exactly.Switch case keys are quoted strings matching the classifier's exact output. "HEALTHY", "AT_RISK", and "UNREACHABLE" must match classify_availability_risk.value verbatim; the defaults: block exists specifically to catch any value that does not match one of those three cases.