OCI Data Flow is a managed Apache Spark service that runs jobs without provisioning clusters, configuring Hadoop, or managing YARN. You submit a job, specify the driver and executor shape, and Data Flow provisions the cluster, runs the job, and terminates everything when done. This post covers Terraform setup, application configuration, private endpoint networking, and scheduling.
Step 1: IAM Policy
resource "oci_identity_policy" "data_flow_policy" {
compartment_id = var.compartment_id
name = "data-flow-policy"
statements = [
"Allow service dataflow to read objects in compartment id COMPARTMENT_OCID",
"Allow service dataflow to manage objects in compartment id COMPARTMENT_OCID where target.bucket.name = 'data-flow-logs'",
"Allow group DataEngineers to manage dataflow-family in compartment id COMPARTMENT_OCID",
"Allow group DataEngineers to use virtual-network-family in compartment id COMPARTMENT_OCID"
]
}
Step 2: Data Flow Application
resource "oci_objectstorage_bucket" "data_flow_logs" {
compartment_id = var.compartment_id
namespace = var.tenancy_namespace
name = "data-flow-logs"
access_type = "NoPublicAccess"
}
resource "oci_dataflow_application" "orders_etl" {
compartment_id = var.compartment_id
display_name = "orders-daily-etl"
description = "Daily ETL from raw to curated zone"
language = "PYTHON"
spark_version = "3.5.0"
file_uri = "oci://data-engineering@NAMESPACE/scripts/orders_etl.py"
archive_uri = "oci://data-engineering@NAMESPACE/deps/orders_etl_deps.zip"
driver_shape = "VM.Standard.E4.Flex"
executor_shape = "VM.Standard.E4.Flex"
num_executors = 4
driver_shape_config { ocpus = 4; memory_in_gbs = 32 }
executor_shape_config { ocpus = 8; memory_in_gbs = 64 }
logs_bucket_uri = "oci://data-flow-logs@NAMESPACE/"
configuration = {
"spark.executor.memoryOverhead" = "4g"
"spark.sql.adaptive.enabled" = "true"
"spark.sql.parquet.compression" = "snappy"
}
subnet_id = var.data_flow_subnet_id
private_endpoint_id = var.data_flow_private_endpoint_id
defined_tags = { "Operations.Environment" = "production", "Operations.ManagedBy" = "terraform" }
}
output "application_id" { value = oci_dataflow_application.orders_etl.id }
Step 3: Run Submission
resource "oci_dataflow_run" "backfill_run" {
compartment_id = var.compartment_id
application_id = oci_dataflow_application.orders_etl.id
display_name = "orders-etl-backfill-20260929"
arguments = [
"--source-bucket", "orders-raw",
"--target-bucket", "data-warehouse",
"--processing-date", "2026-09-28"
]
logs_bucket_uri = "oci://data-flow-logs@NAMESPACE/runs/20260929/"
}
Step 4: Run Status Alarm
resource "oci_monitoring_alarm" "data_flow_job_failed" {
compartment_id = var.compartment_id
display_name = "data-flow-job-failed"
is_enabled = true
metric_compartment_id = var.compartment_id
namespace = "oci_dataflow"
query = "RunsFailed[1h]{applicationId = 'DATAFLOW_APP_OCID'}.sum() > 0"
severity = "CRITICAL"
pending_duration = "PT5M"
destinations = [var.data_team_topic_id]
body = "OCI Data Flow ETL job failed. Check run logs in the data-flow-logs bucket for stack trace and error details."
}
Operational Notes
Data Flow charges for the duration the Spark cluster runs including startup time. For jobs that complete in under 5 minutes, cluster startup overhead is proportionally significant. Batch larger work into longer-running jobs rather than many small ones to keep cost per unit of work low. Unlike persistent clusters, there is no idle time cost because the cluster terminates when the job ends.
Use a private endpoint to keep Data Flow traffic inside your VCN. The private endpoint connects the Data Flow managed runtime to your Object Storage through private IP, avoiding the public Object Storage endpoint. Without a private endpoint, any bucket-level network access controls that restrict to private CIDRs will block Data Flow runs.
Regards,
Osama
#OCI #OracleCloud #DataFlow #ApacheSpark #ETL #Terraform #IaC #TechBlog #Oracle #DataEngineering #Analytics #Serverless #BigData #PySpark #OracleCloudInfrastructure #ObjectStorage #DataPipeline #CloudData #DataPlatform #DataWarehouse
Leave a comment