okustera_workflow_dag (Resource)
Manages an Apache Airflow DAG workflow definition in Okustera. DAG scripts are stored directly in tenant Ceph RADOS Gateway S3 buckets with real-time bidirectional synchronization to the Airflow KubernetesExecutor scheduler cluster.
Example Usage
Automated Customer ETL Pipeline
resource "okustera_workflow_dag" "data_pipeline" {
dag_id = "tenant_customer_etl"
filename = "customer_etl.py"
is_paused = false
content = <<-EOT
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta
default_args = {
'owner': 'data-engineering',
'depends_on_past': False,
'email_on_failure': False,
'retries': 2,
'retry_delay': timedelta(minutes=5),
}
with DAG(
dag_id='tenant_customer_etl',
default_args=default_args,
description='Automated customer synchronization ETL pipeline',
schedule_interval='@daily',
start_date=datetime(2026, 1, 1),
catchup=False,
tags=['etl', 'customers', 'analytics'],
) as dag:
extract = BashOperator(
task_id='extract_records',
bash_command='echo "Extracting customer data..."',
)
transform = BashOperator(
task_id='transform_records',
bash_command='echo "Transforming and validating schema..."',
)
load = BashOperator(
task_id='load_warehouse',
bash_command='echo "Loading transformed records into analytics warehouse..."',
)
extract >> transform >> load
EOT
}
output "workflow_dag_id" {
value = okustera_workflow_dag.data_pipeline.dag_id
}
output "workflow_schedule" {
value = okustera_workflow_dag.data_pipeline.schedule_interval
}
Schema
Required
content(String) Python workflow DAG source code content.dag_id(String, Forces new resource) Unique DAG ID identifier referenced in the Python script.
Optional
filename(String) Target file name stored in Ceph S3 (defaults to<dag_id>.py).is_paused(Boolean) Whether the workflow scheduling is paused (defaults tofalse).
Read-Only
description(String) Human-readable description parsed from DAG definition.fileloc(String) Storage path inside the Airflow scheduler container.id(String) Unique workflow DAG identifier.is_active(Boolean) Whether the DAG is recognized as active by the Airflow scheduler.last_parsed_time(String) Timestamp when Airflow parsed the DAG.last_run_start_date(String) Start timestamp of the most recent DAG run.last_run_state(String) State of the most recent DAG run.next_dagrun(String) Estimated timestamp for the next automated execution run.owners(List of String) List of owners responsible for this pipeline.schedule_interval(String) Schedule interval expression (cron or timetable alias).tags(List of String) Organizational tags attached to this workflow.timetable_description(String) Human-readable timetable schedule description.
Import
Import is supported using the unique dag_id:
terraform import okustera_workflow_dag.data_pipeline tenant_customer_etl