apache / apache/iggy

feat(connectors): add Apache Airflow trigger connector

Open
#3,715 2 comments 0 reactions 1 assignee Claimed by @avirajkhare00 View on GitHub
connectors good first issue
Dominant language
Rust
Stars
4.9k
Forks
432
Avg merge
2d 10h
Merged PRs (30d)
173

Description

### Description

Add an Apache Airflow connector so Iggy can drive event-driven DAGs via the Airflow REST API.

This is tracked on the connector ecosystem roadmap (#2753) under **Workflow & Orchestration**:

| Target | Type | Priority | Comp | Suggested stack |
|--------|------|:--------:|:----:|-----------------|
| [Apache Airflow](https://airflow.apache.org/) | Trigger | P4 | 1/4 | `reqwest` (REST API) |

Roadmap intent: **Iggy sensor/trigger for event-driven DAGs**. There is no Airflow sink/source (or trigger plugin) in-tree today, no dedicated implementation issue before this one, and no open PR claiming the work.

**Motivation**

- Airflow is the industry-standard workflow orchestrator; pairing it with Iggy enables stream-driven DAG triggers instead of polling-only sensors.
- Existing connectors cover DBs, search, lakehouse, and HTTP egress, but nothing in the workflow/orchestration category is implemented yet (Airflow, Temporal, Prefect, Dagster are all still roadmap-only).
- A first cut can likely reuse patterns from `http_sink` (auth, retries, batching) while specializing for Airflow's trigger/DAG-run APIs.

**Suggested v1 scope (trigger)**

- Consume messages from configured Iggy streams/topics
- Trigger DAG runs via Airflow REST API (e.g. create DAG run / trigger endpoint)
- Config: Airflow base URL, auth (token/basic), `dag_id`, optional conf/payload mapping from message body, retry policy
- Map transient HTTP errors (5xx, timeouts) to retry; map permanent client errors (auth, missing DAG) to non-retry
- `open()` connectivity check against Airflow health/version endpoint; structured logging of trigger outcomes
- Unit tests + example config under `core/connectors/runtime/example_config/connectors/`

Out of scope for v1 (unless someone wants to expand): full Airflow **source** (task/DAG state → Iggy), Prefect/Dagster parity, Airflow provider package on the Airflow side.

### Affected area / component

Connectors

### Proposed solution

New sink-style (or dedicated trigger) plugin, e.g. `airflow_sink` / `airflow_trigger`, under `core/connectors/sinks/`, built on `iggy_connector_sdk` and `reqwest`, following existing sink lifecycle (`open` / `consume` / `close`) and config conventions.

### Alternatives considered

- **`http_sink` only**: works for a one-off webhook but does not encode Airflow auth, DAG-run semantics, conf mapping, or operator-friendly defaults.
- **Airflow-side provider/sensor only**: useful later, but does not give Iggy-native plugin packaging, runtime metrics, or the same ops model as other connectors.
- **Wait for a generic "workflow trigger" abstraction**: delays a concrete Airflow integration that the roadmap already calls out.

### Contribution

- [x] I'm willing to submit a pull request to implement this feature

### Good first issue

- [x] I think this could be a good first issue for a new contributor

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.