New to Kestra?
Use blueprints to kickstart your first workflows.
Audit DuckDB index coverage and database file size with its own catalog functions, alert on risk, and optionally create the missing index.
A table queried by a column with no backing index degrades the same way in DuckDB
as in any SQL engine: scans get slower as the table grows, without ever throwing an
error that points at the real cause. A DuckDB file that keeps growing without
anyone checking pragma_database_size() is the other half of the same blind spot -
both are easy to see with one query and easy to never check otherwise. This
blueprint reads DuckDB's own catalog functions - duckdb_tables() and
duckdb_indexes() for index coverage, pragma_database_size() for file size -
classifies the result into a three-state risk level, alerts on anything but a
clean bill of health, and can create a missing index with a plain CREATE INDEX
statement, behind two independent gates.
periodic_governance_check (io.kestra.plugin.core.trigger.Schedule) runs every 6 hours, always forcing auto_remediate: "false" and dry_run: "true" through its own inputs: override regardless of this flow's defaults. Shipped disabled so you can validate a manual run first.check_table_and_index (io.kestra.plugin.jdbc.duckdb.Query, fetchType: FETCH_ONE) joins duckdb_tables() against a correlated count from duckdb_indexes(), populating row.table_name, row.estimated_size, row.index_count, row.matching_indexes - or no row at all if target_table does not exist.check_storage_size (io.kestra.plugin.jdbc.duckdb.Query, fetchType: FETCH_ONE) reads pragma_database_size(), populating row.database_size, row.block_size, row.used_blocks, row.total_blocks.evaluate_table_presence (io.kestra.plugin.core.flow.If) branches on whether check_table_and_index returned a row.else (UNREACHABLE): alert_table_unreachable reports that target_table does not exist in this database.then (table found): nested evaluate_governance_risk (io.kestra.plugin.core.flow.If) branches on matching_indexes == 0 or used_blocks > max_used_blocks.else (HEALTHY): log_governance_healthy records a clean check.then (AT_RISK): route_remediation (io.kestra.plugin.core.flow.Switch on auto_remediate) either considers remediation ("true", further gated by check_has_missing_index - only a missing index can be created, not storage reclaimed - and then by dry_run through check_dry_run) or only logs that remediation is disabled ("false"). When every gate allows it, create_missing_index runs CREATE INDEX ... ON target_table(required_indexed_column). alert_governance_risk always fires to Slack regardless of which remediation path ran.log_audit_complete always runs last, printing the headline numbers regardless of which branch fired.errors block alerts Slack separately if the flow itself fails outright - a locked or corrupted database file must not read as "no risk".CREATE INDEX statement on the one table and column you named - it never touches row data or any other table.auto_remediate, dry_run) between "gap detected" and "index actually created".dbt-model-runtime-regression-alert.yaml writes to) who need an early signal when a hot query column was never indexed.Query driving catalog introspection (duckdb_tables(), duckdb_indexes(), pragma_database_size()) and a DDL statement (CREATE INDEX), not just file-based SQL analytics.duckdb_tables() and pragma_database_size() run by hand answer one moment in time; they do not decide a cadence, classify the result into a three-state risk level, gate a schema change behind two independent confirmations, or notify anyone. Kestra supplies the schedule, the HEALTHY/AT_RISK/UNREACHABLE classification as a first-class branch, two independent authorization gates before any index is created, and an execution history that shows exactly when a table went uncovered or a file outgrew its budget.
target_table already created in it.Local testing:
This flow has no Docker image of its own to run, since DuckDB is embedded directly inside the Kestra Worker process via its JDBC driver, not a separate server. docker network connect YOUR_KESTRA_NETWORK <container> and a host-port mapping away from 8080 therefore do not apply here - there is no container to connect. Instead, create the database file and table once, from any environment with a DuckDB client, for example: docker run --rm -v /tmp/kestra-duckdb:/data python:3.12-slim bash -c "pip install duckdb -q && python3 -c \"import duckdb; con = duckdb.connect('/data/analytics.duckdb'); con.execute('CREATE TABLE IF NOT EXISTS events (event_id VARCHAR, payload VARCHAR)')\"". Ensure database_path in this flow resolves to the same file from the Kestra Worker's own filesystem (mount the same host path into the Worker container if Kestra itself runs in Docker).
SLACK_WEBHOOK_URL: Slack incoming webhook used by alert_governance_risk, alert_table_unreachable, and the errors block.{{ secret(...) }} for check_table_and_index, check_storage_size, or create_missing_index.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 the webhook 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.database_path (STRING, default /tmp/kestra-duckdb/analytics.duckdb): file path rendered into every task's url.target_table (STRING, default events): table audited and, if remediated, altered with a new index.required_indexed_column (STRING, default event_id): column expected to be index-covered.max_used_blocks (INT, default 100000): storage-growth ceiling on pragma_database_size().used_blocks.auto_remediate (SELECT: "false", "true"; default "false"): must be "true" for remediation to even be considered.dry_run (SELECT: "true", "false"; default "true"): must be explicitly "false", together with auto_remediate: "true", for CREATE INDEX to actually run.outputs.check_table_and_index.row.table_name / .estimated_size / .index_count / .column_count / .matching_indexes: catalog facts for target_table, every run.outputs.check_storage_size.row.database_size / .block_size / .total_blocks / .used_blocks / .free_blocks: file-size facts for the database, every run.events table with no index on event_id, using the command in Prerequisites (or any DuckDB client pointed at the same file).SLACK_WEBHOOK_URL secret.alert_governance_risk fires (AT_RISK: the table exists, no index covers event_id yet).auto_remediate: "true" and dry_run: "true" and confirm log_dry_run_remediation describes the index without creating it.auto_remediate: "true" and dry_run: "false" and confirm create_missing_index runs, then verify with SELECT * FROM duckdb_indexes() that the index now exists.target_table at a table name that does not exist and re-run to confirm alert_table_unreachable fires (UNREACHABLE) instead of either HEALTHY or AT_RISK.periodic_governance_check once you trust the check; it always runs in the safe auto_remediate: "false" / dry_run: "true" mode regardless of what you leave the flow's own defaults set to.target_table/required_indexed_column pairs with io.kestra.plugin.core.flow.Loop to run the same gate across every governed table in one database file.VACUUM or CHECKPOINT as a separate, more cautiously gated remediation path for the storage-growth half of AT_RISK, since CREATE INDEX intentionally does not address it.check_storage_size's query by database_name instead of LIMIT 1 if this connection ever attaches more than one database.has_primary_key from duckdb_tables() to the HEALTHY/AT_RISK decision for tables where a primary key, not just a secondary index, is the real governance requirement.Query task opens the same local file directly through the JDBC driver running inside the Kestra Worker. If the Worker runs in a container, database_path must resolve to a file visible inside that container, not just on the host.pragma_database_size() LIMIT 1 assumes a single attached database. If more than one database is ever attached to this connection (via ATTACH), the row returned may not be the one you expect - filter by database_name explicitly in that case, per How to extend.CREATE INDEX here has no IF NOT EXISTS guard. Re-running create_missing_index after the index already exists will error on the duplicate name; check_has_missing_index prevents this in the normal gated path by only firing when matching_indexes == 0, but a manually re-run task outside that guard is not protected.expressions::VARCHAR LIKE check is a substring match, not an exact column match. A column named event_id will also match an index on a differently named column that happens to contain that substring (for example source_event_id) - acceptable for a compliance gate erring toward false negatives on "no index", but worth tightening with an exact expression match if your schema has overlapping column names.CREATE INDEX DDL go through the same generic io.kestra.plugin.jdbc.duckdb.Query task used elsewhere in this repo's own DuckDB flows.auto_remediate and dry_run are independent gates, not a single boolean. Both must be "true"/"false" respectively for CREATE INDEX to run; setting only one leaves the other still blocking.Switch case keys "true"/"false" are quoted strings. auto_remediate renders as the literal string "true" or "false"; unquoted true:/false: YAML map keys would parse as booleans instead and would not match.duckdb_tables(), duckdb_indexes(), and pragma_database_size() column names are taken verbatim from DuckDB's own official documentation, and the Query task's properties (url, sql, fetchType) are taken verbatim from this repo's own existing DuckDB flows - but no Kestra engine was run to execute this flow end to end.