feat(connectors): add Apache Fluss sink connector
- Dominant language
- Rust
- Stars
- 4.9k
- Forks
- 432
- Avg merge
- 2d 10h
- Merged PRs (30d)
- 174
Description
### Description
Add a **sink connector** that consumes messages from Iggy topics and writes them into [Apache Fluss](https://fluss.apache.org/) tables.
Fluss provides durable, offset-ordered **log tables** and **primary-key tables** with sub-millisecond lookups and lakehouse tiering. The official Rust client ([`fluss-rs`](https://crates.io/crates/fluss-rs)) speaks **Apache Arrow** natively, which pairs well with JSON/Avro payloads on the Iggy side and Arrow-oriented egress (for example Iceberg).
Use cases:
- Land Iggy stream data into Fluss as the real-time analytics / serving layer.
- Feed Fluss hot tier from Iggy connectors pipelines (sources → transforms → Fluss).
- Use Fluss PK tables for materialized views keyed by business identifiers from Iggy events.
## Affected area / component
- Connectors (new `fluss_sink` plugin under `core/connectors/sinks/`)
- Connectors runtime configuration examples
- Integration tests (optional follow-up; may need Fluss test fixture)
### Affected area / component
_No response_
### Proposed solution
Implement `iggy_connector_fluss_sink` as a `cdylib` plugin implementing `iggy_connector_sdk::Sink`, following conventions from `postgres_sink` and `iceberg_sink`.
### Client dependency
Use [`fluss-rs`](https://github.com/apache/fluss-rust). Pin a released version in the workspace `Cargo.toml` when adding the crate.
### Write modes (phased)
| Phase | Fluss table type | Write API | Iggy input |
| ----- | ---------------- | --------- | ---------- |
| 1 | Log table | `AppendWriter` / Arrow `RecordBatch` append | `Schema::Json` (row objects) |
| 1b | Log table | Arrow IPC `Schema::Raw` passthrough | Zero-copy when upstream sends Arrow batches |
| 2 | Primary-key table | `UpsertWriter` upsert/delete | JSON with PK columns + optional `op` header |
### Message mapping
- **JSON payloads**: map object fields to `GenericRow` using Fluss table schema (fail batch on type mismatch; log and skip or dead-letter per existing sink patterns).
- **Avro payloads**: optional transform via existing `avro_convert` before write.
- **Arrow IPC (`Schema::Raw`)**: decode to `RecordBatch` and call Fluss batch append API when `payload_format = "arrow_ipc"`.
### Routing
Phase 1: single destination `(database, table)` from config.
Phase 2 (optional): static or dynamic routing by message field (same pattern as `iceberg_sink` `dynamic_route_field`).
### Example configuration
```toml
[plugin_config]
bootstrap_servers = "localhost:9123"
database = "mydb"
table = "events"
table_type = "log" # "log" | "primary_key"
write_mode = "append" # "append" | "upsert" | "delete"
payload_format = "json" # "json" | "arrow_ipc"
batch_size = 500
flush_interval = "1s"
# primary_key_columns = ["id"] # required when table_type = "primary_key"
# dynamic_routing = false
# dynamic_route_field = "fluss_table" # "database.table"
```
### Files / layout (expected)
```text
core/connectors/sinks/fluss_sink/
├── Cargo.toml
├── config.toml
├── README.md
└── src/
├── lib.rs
├── config.rs
└── writer.rs
```
Register in workspace `Cargo.toml`, `core/connectors/sinks/README.md`, and add a runtime example under `core/connectors/runtime/example_config/connectors/`.
### Example plugins
- Sink lifecycle / batch consume: `postgres_sink`, `mongodb_sink`
- JSON → typed row mapping: `postgres_sink`
- Arrow batch path: `iceberg_sink` (`arrow_json` / `RecordBatch`)
## Acceptance criteria
- [ ] `fluss_sink` plugin builds as `cdylib` and loads in connectors runtime
- [ ] Consumes Iggy topic batches and appends to configured Fluss log table
- [ ] JSON payload mapping covered by unit tests; integration test or documented manual test plan against local Fluss cluster
- [ ] Transient Fluss errors map to retryable `Error` variants; permanent schema errors fail fast
- [ ] Config documented in plugin `README.md` with minimal runnable example
- [ ] `cargo fmt`, `cargo sort --no-format`, `cargo clippy`, `cargo test -p iggy_connector_fluss_sink` pass
## References
- Apache Fluss: https://fluss.apache.org/
- Fluss Rust client: https://github.com/apache/fluss-rust
- `fluss-rs` on crates.io: https://crates.io/crates/fluss-rs
- Fluss Arrow / RecordBatch writes: https://fluss.apache.org/blog/fluss_rust_client_release/
- Related: source connector issue (file separately; link when available)
### Alternatives considered
## Alternatives considered
1. **Write via Flink Table API only** — excludes Iggy-native connector users.
2. **Iceberg sink only** — Fluss hot tier + PK lookup is a different latency/consistency profile than object-store lake tables.
3. **Custom Protobuf client** — duplicate of maintained `fluss-rs`; avoid.
### Contribution
- [ ] I'm willing to submit a pull request to implement this feature
### Good first issue
- [ ] I think this could be a good first issue for a new contributor
Contributor guide
Research direction
Start by reading the postgres_sink and iceberg_sink implementations, then inspect workspace Cargo.toml and the connectors runtime example configuration. Use the expected core/connectors/sinks/fluss_sink layout and implement the phased scope described, beginning with JSON appends to a configured log table. Done means the cdylib loads, unit tests cover JSON mapping, configuration is documented, and cargo test -p iggy_connector_fluss_sink passes.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- rust
- Domain
- backend, data-engineering
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Quiet
- Clarity
- Mostly clear
- Newbie friendliness
- 38/100