Airflow
Move data between Airflow DAGs and RustFS with the Amazon S3 provider.
This guide connects Apache Airflow — the workflow orchestration platform — to RustFS through the Amazon S3 provider's hooks, operators, and sensors. You will register a custom-endpoint connection, run a DAG that writes an object to a RustFS bucket, waits for a key with S3KeySensor, and reads the object back with S3Hook. The workflow was verified with apache/airflow:3.3.2 (standalone, SequentialExecutor) and the apache-airflow-providers-amazon provider against rustfs/rustfs-x86-musl:v2.3.1.
You need Docker, or an existing Airflow installation. This deployment is intended for local integration testing, not production.
Architecture
Every S3 interaction inside a DAG goes through the provider's S3 client, pointed at RustFS by the connection's endpoint_url. Operators, sensors, and hooks share the same connection object.
Run Airflow
Start a standalone instance with examples disabled, and create the demo bucket:
docker run -d --name airflow --network oo-rustfs_default -p 8080:8080 \
-e AIRFLOW__CORE__LOAD_EXAMPLES=False \
-v "$PWD/dags":/opt/airflow/dags \
apache/airflow:3.3.2 standalone
rc mb rustfs/airflow-demoThe image ships with all providers preinstalled, including apache-airflow-providers-amazon.
Register the RustFS connection
The S3 provider reads its endpoint from the connection's extra field. Replace all connection placeholders:
docker exec airflow airflow connections add rustfs \
--conn-type aws \
--conn-extra '{"endpoint_url": "http://<your-rustfs-endpoint>:9000", "region_name": "us-east-1", "aws_access_key_id": "<your-access-key>", "aws_secret_access_key": "<your-secret-key>"}'The keys aws_access_key_id and aws_secret_access_key inside --conn-extra supply credentials; endpoint_url redirects the boto3 client from AWS to RustFS.
Write the DAG
The DAG writes an object with an operator, waits for the key with a sensor, and reads it back with the hook:
import datetime
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
from airflow.providers.amazon.aws.operators.s3 import S3CreateObjectOperator
from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor
from airflow.sdk import dag, task
@dag(
schedule=None,
start_date=datetime.datetime(2026, 1, 1),
catchup=False,
tags=["rustfs"],
)
def rustfs_demo():
create = S3CreateObjectOperator(
task_id="write_object",
s3_bucket="airflow-demo",
s3_key="dags/airflow-put.txt",
data="written by airflow to rustfs",
aws_conn_id="rustfs",
replace=True,
)
wait = S3KeySensor(
task_id="wait_for_object",
bucket_key="dags/airflow-put.txt",
bucket_name="airflow-demo",
aws_conn_id="rustfs",
timeout=120,
poke_interval=10,
mode="reschedule",
)
@task
def read_object():
hook = S3Hook(aws_conn_id="rustfs")
body = hook.read_key(key="dags/airflow-put.txt", bucket_name="airflow-demo")
print("read back:", body)
assert body == "written by airflow to rustfs"
create >> [wait, read_object()]
rustfs_demo()Note the import paths: S3CreateObjectOperator lives in the operators module while S3KeySensor lives in the sensors module — importing both from one place fails.
Unpause and trigger
New DAGs start paused, and a trigger fired while paused stays queued forever. Unpause first, then trigger:
docker exec airflow airflow dags unpause rustfs_demo
docker exec airflow airflow dags trigger rustfs_demoWatch the run finish:
docker exec airflow airflow dags list-runs rustfs_demo | head -3dag_id run_id state
rustfs_demo manual__2026-09-29T13:15:57.332688+00:00 successAll three tasks succeed: write_object, wait_for_object, and read_object.
Verify objects in RustFS
List the bucket prefix:
rc ls rustfs/airflow-demo/ -r
rc cat rustfs/airflow-demo/dags/airflow-put.txt[2026-09-29 13:16:01] 28 B dags/airflow-put.txt
written by airflow to rustfs
Stop or reset
To tear down the demo while keeping the bucket objects:
docker rm -f airflowTo delete the stored data:
rc rm rustfs/airflow-demo/ --recursive --forceTroubleshooting
Dag runs stay queued after triggering
The DAG is paused. New DAGs are paused by default in Airflow 3, and runs triggered in that state never execute. Run airflow dags unpause rustfs_demo; queued runs then start on their own.
cannot import name 'S3KeySensor' from 'airflow.providers.amazon.aws.operators.s3'
The sensor lives in a separate module: from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor.
Tasks fail with connection errors
The endpoint_url must be reachable from the Airflow container — use the Docker network hostname for RustFS, not localhost. Airflow 3 serves its health endpoint under /api/v2/monitor/health if you need to check component status.
Next steps
- Review S3 compatibility notes before adopting additional provider hooks.
- Create dedicated production credentials with Access Key Management.
- Follow the Amazon provider documentation for transfer operators such as
S3ToLocalFilesystemOperatorandLocalFilesystemToS3Operator.