Skip to main content

Workflow & Data Orchestration (Apache Airflow)

Phase 11 Deployed & Production Ready

Managed Apache Airflow PaaS is fully deployed and operational on the Okustera platform in namespace airflow-system. Workflows are served via Apache APISIX at https://airflow.okustera.com, backed by CloudNativePG PostgreSQL 16, on-premise Ceph RADOS Gateway S3 with OpenStack Barbican KMS encryption, and the serverless KubernetesExecutor.

Okustera Managed Apache Airflow delivers enterprise-grade workflow orchestration and ETL data pipeline scheduling as a zero-license-tax sovereign cloud service—the open-source equivalent to AWS MWAA and Google Cloud Composer.

Tenants can provision dedicated, autoscaling Apache Airflow environments with a single click in the Okustera Cloud Portal or declaratively using Terraform, completely eliminating the operational overhead of managing Airflow schedulers, metadata databases, and task workers.


Architecture: Cloud-Native Serverless Elasticity​

Rather than running static, always-on Celery worker VMs that consume idle compute, Okustera Airflow utilizes the KubernetesExecutor on top of the tenant Kubernetes workload plane. Every individual task run in a Directed Acyclic Graph (DAG) executes in a dedicated, ephemeral Kubernetes pod that scales up on demand and terminates upon completion.


Key Capabilities​

  • Serverless Task Elasticity: The KubernetesExecutor provisions on-demand task pods with custom CPU, memory, and GPU resource constraints, scaling from 0 to hundreds of parallel tasks instantly.
  • 100% Sovereign Backing: Backed natively by CloudNativePG (PostgreSQL 16) for metadata persistence and Ceph RADOS Gateway (RGW S3) for DAG and log storage—zero proprietary external cloud dependencies.
  • Dual DAG Synchronization:
    • Git-Sync Sidecar: Continuously synchronizes DAG repositories from GitHub, GitLab, or internal Git servers.
    • Direct S3 Upload: Drag-and-drop Python DAG scripts directly into s3://omc-airflow-dags/<tenant-id>/ from the Okustera Cloud Portal UI.
  • Unified Security & Governance: Fully integrated with OpenStack Keystone multi-tenant IAM and Apache APISIX for TLS termination and token authentication.
  • Centralized Observability: Automatic streaming of task logs to Ceph S3, metrics to Prometheus/Grafana, and audit logs to Grafana Loki.

Declarative Management via Terraform​

Workflows, pipelines, and DAG definitions are managed declaratively using the official Okustera Terraform Provider:

# Inspect operational health and telemetry of the Airflow cluster
data "okustera_workflow_status" "cluster" {}

# Deploy a production ETL DAG pipeline with S3 synchronization
resource "okustera_workflow_dag" "analytics_pipeline" {
dag_id = "tenant_analytics_pipeline"
filename = "analytics_pipeline.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,
'retries': 2,
'retry_delay': timedelta(minutes=5),
}

with DAG(
dag_id='tenant_analytics_pipeline',
default_args=default_args,
description='Nightly customer analytics and financial billing DAG',
schedule_interval='@daily',
start_date=datetime(2026, 1, 1),
catchup=False,
tags=['production', 'analytics', 'finops'],
) as dag:

task_extract = BashOperator(
task_id='extract_metrics',
bash_command='echo "Extracting telemetry metrics..."',
)

task_process = BashOperator(
task_id='process_aggregates',
bash_command='echo "Processing data aggregates..."',
)

task_extract >> task_process
EOT
}

output "workflow_status" {
value = data.okustera_workflow_status.cluster.status
}

output "dag_parsed_schedule" {
value = okustera_workflow_dag.analytics_pipeline.schedule_interval
}

Refer to the okustera_workflow_dag Resource Documentation and okustera_workflow_status Data Source Documentation for complete schema specifications.


Example: Sovereign S3-to-PostgreSQL ETL DAG​

The following Python script illustrates an Apache Airflow DAG executing on Okustera, extracting customer event files from Ceph RGW S3, processing records in parallel with pandas, and loading transformed tables into CloudNativePG PostgreSQL:

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
import boto3
import psycopg2
import pandas as pd
import io

default_args = {
"owner": "data-team",
"depends_on_past": False,
"email_on_failure": False,
"retries": 2,
"retry_delay": timedelta(minutes=5),
}

def extract_and_transform_s3_events():
"""Extract telemetry data from Ceph S3 and transform with pandas."""
s3 = boto3.client(
"s3",
endpoint_url="https://s3.okustera.com",
aws_access_key_id="YOUR_OKUSTERA_ACCESS_KEY",
aws_secret_access_key="YOUR_OKUSTERA_SECRET_KEY",
)

# Download raw batch file
obj = s3.get_object(Bucket="customer-events", Key="batches/latest.csv")
df = pd.read_csv(io.BytesIO(obj["Body"].read()))

# Clean and aggregate metrics
df_clean = df.dropna(subset=["event_id", "timestamp"])
df_clean["processed_at"] = datetime.utcnow().isoformat()

# Save transformed intermediate dataset back to S3
csv_buffer = io.StringIO()
df_clean.to_csv(csv_buffer, index=False)
s3.put_object(
Bucket="customer-events",
Key="processed/latest_cleaned.csv",
Body=csv_buffer.getvalue(),
)
print(f"Successfully processed {len(df_clean)} records.")

def load_to_postgresql():
"""Load transformed dataset into CloudNativePG PostgreSQL database."""
conn = psycopg2.connect(
host="production-crm-db-rw.default.svc.cluster.local",
port=5432,
dbname="analytics",
user="postgres",
password="<PASSWORD>",
)
cursor = conn.cursor()
cursor.execute("""
CREATE TABLE IF NOT EXISTS customer_events_summary (
event_id VARCHAR(64) PRIMARY KEY,
processed_at TIMESTAMP,
event_count INT
);
""")
conn.commit()
cursor.close()
conn.close()
print("Database sync completed successfully.")

with DAG(
"okustera_s3_postgres_pipeline",
default_args=default_args,
description="Sovereign ETL data pipeline running on KubernetesExecutor",
schedule_interval="0 2 * * *", # Daily at 2:00 AM UTC
start_date=datetime(2026, 1, 1),
catchup=False,
tags=["sovereign", "etl", "s3", "postgres"],
) as dag:

t1 = PythonOperator(
task_id="extract_transform_s3",
python_callable=extract_and_transform_s3_events,
)

t2 = PythonOperator(
task_id="load_postgresql",
python_callable=load_to_postgresql,
)

t1 >> t2

Integration with the Okustera Cloud Portal​

Tenants manage Airflow pipelines directly from the Workflows console in the Okustera Cloud Portal:

  • Workspaces Overview: Real-time status pills (HEALTHY, UPDATING, DEGRADED), scheduler uptime, and active pod resource utilization.
  • DAG Catalog & Run History: Live matrix of loaded DAGs, active/failed runs, and one-click manual execution triggers.
  • Integrated S3 Log Streamer: Direct terminal viewing of task execution logs fetched on-the-fly from s3://omc-airflow-logs.

Managing Workflows in the Cloud Portal​

The Okustera Cloud Portal (/workflows) provides an intuitive web interface for data engineers and pipeline operators:

1. Uploading a Python DAG​

  1. Navigate to Workflows in the left sidebar.
  2. Click Upload DAG:
    • Drag and drop your .py DAG file into the upload zone (or click to browse).
  3. The portal uploads the script directly to your tenant's S3 DAG prefix (s3://omc-airflow-dags/<tenant_id>/).
  4. The Airflow Scheduler synchronizes the new file within 10 seconds and registers the DAG in your catalog.

2. Triggering and Inspecting DAG Runs​

  1. In the DAG list, locate your workflow (e.g. tenant_analytics_pipeline).
  2. Click Trigger DAG:
    • Optionally provide custom runtime configuration JSON (e.g., {"date": "2026-09-28", "batch_size": 500}).
  3. Click Confirm Trigger.
  4. The execution timeline opens automatically, displaying the status of spawned KubernetesExecutor task pods.

3. Viewing Live Task Execution Logs​

  1. Click on any running or completed task in the execution graph.
  2. Select Task Logs.
  3. Real-time task execution logs stream directly into the portal from Ceph S3, eliminating the need to log into Kubernetes nodes or run kubectl logs.

4. Pausing and Deleting DAGs​

  • Toggle the Pause / Unpause switch to halt scheduled cron triggers while preserving execution history.
  • Click Delete DAG to remove the script from S3 and clean up scheduler metadata.

REST API Reference​

  • GET /api/v1/workflows/status — Inspect operational health of Airflow Scheduler and executor pods.
  • GET /api/v1/workflows/dags — List workflow DAGs accessible to current tenant.
  • GET /api/v1/workflows/dags/{dag_id} — Retrieve metadata, schedule, and execution timetable for a DAG.
  • PATCH /api/v1/workflows/dags/{dag_id} — Pause or unpause a workflow DAG.
  • POST /api/v1/workflows/dags/{dag_id}/trigger — Trigger immediate DAG execution with optional JSON parameters.
  • GET /api/v1/workflows/dags/{dag_id}/runs — List historical and active DAG runs with state and duration.
  • GET /api/v1/workflows/dags/{dag_id}/runs/{run_id}/tasks — List task instances and execution status for a run.
  • GET /api/v1/workflows/dags/{dag_id}/runs/{run_id}/tasks/{task_id}/logs — Retrieve streaming task execution logs.
  • POST /api/v1/workflows/dags/upload — Upload Python DAG file to S3 with automatic scheduler synchronization.
  • DELETE /api/v1/workflows/dags/{dag_id} — Delete DAG file from storage and purge scheduler metadata.