Skip to content

add apache Fluss source connector #3688

Description

@ryerraguntla

Description

Add a source connector that reads from Apache Fluss tables and publishes records into Iggy streams/topics.

Fluss is streaming storage for real-time analytics and lakehouse hot tiers. Its official Rust client (fluss-rs) uses Apache Arrow (RecordBatch / GenericRow) end-to-end. That aligns with Iggy's connectors subsystem and with Arrow-friendly downstream sinks (for example Iceberg).

Use cases:

  • Bridge Fluss log-table changelogs into Iggy for fan-out to multiple consumers or SDKs.
  • Ingest Fluss primary-key table CDC-style changes into Iggy for event-driven apps.
  • Unify Fluss hot-tier data with Iggy-based routing, transforms, and connector egress.

Affected area / component

Connectors

Proposed solution

Implement iggy_connector_fluss_source as a cdylib plugin implementing iggy_connector_sdk::Source, following conventions from postgres_source and influxdb_source.

Client dependency

Use fluss-rs (Tokio async, Arrow-native). Pin a released version in the workspace Cargo.toml when adding the crate.

Table modes (phased)

Phase Fluss table type Read API Notes
1 Log table subscribe + poll / batch scanner Offset-ordered append-only; primary use case
2 Primary-key table Changelog / scan APIs as exposed by fluss-rs Upsert/delete events; confirm API stability in 0.1.x

Message mapping

Default: decode Fluss rows to Schema::Json via existing connector JSON encoders (one object per row, column names from Fluss schema).

Optional (config flag): emit Schema::Raw Arrow IPC bytes per RecordBatch for zero-copy analytics pipelines (consumers decode Arrow client-side).

Offset / state

Persist Fluss log offsets (and partition/bucket identifiers if applicable) in connector state via the runtime state API (msgpack), same pattern as other sources. On restart, resume from last committed offset.

Draft configuration

[plugin_config]
bootstrap_servers = "localhost:9123"
database = "mydb"
table = "events"
table_type = "log"          # "log" | "primary_key"
starting_offset = "earliest" # "earliest" | "latest" | explicit offset
batch_size = 1000
poll_interval = "500ms"
payload_format = "json"       # "json" | "arrow_ipc"
# Optional auth when supported by cluster
# sasl_mechanism = "PLAIN"
# sasl_username = "..."
# sasl_password = "..."

Implement iggy_connector_fluss_source as a cdylib plugin implementing iggy_connector_sdk::Source, following conventions from postgres_source and influxdb_source.

Client dependency

Use fluss-rs (Tokio async, Arrow-native). Pin a released version in the workspace Cargo.toml when adding the crate.

Table modes (phased)

Phase Fluss table type Read API Notes
1 Log table subscribe + poll / batch scanner Offset-ordered append-only; primary use case
2 Primary-key table Changelog / scan APIs as exposed by fluss-rs Upsert/delete events; confirm API stability in 0.1.x

Message mapping

Default: decode Fluss rows to Schema::Json via existing connector JSON encoders (one object per row, column names from Fluss schema).

Optional (config flag): emit Schema::Raw Arrow IPC bytes per RecordBatch for zero-copy analytics pipelines (consumers decode Arrow client-side).

Offset / state

Persist Fluss log offsets (and partition/bucket identifiers if applicable) in connector state via the runtime state API (msgpack), same pattern as other sources. On restart, resume from last committed offset.

Draft configuration

[plugin_config]
bootstrap_servers = "localhost:9123"
database = "mydb"
table = "events"
table_type = "log"          # "log" | "primary_key"
starting_offset = "earliest" # "earliest" | "latest" | explicit offset
batch_size = 1000
poll_interval = "500ms"
payload_format = "json"       # "json" | "arrow_ipc"
# Optional auth when supported by cluster
# sasl_mechanism = "PLAIN"
# sasl_username = "..."
# sasl_password = "..."

Files / layout (expected)

core/connectors/sources/fluss_source/
├── Cargo.toml
├── config.toml
├── README.md
└── src/lib.rs

Register in workspace Cargo.toml, core/connectors/sources/README.md, and add a runtime example under core/connectors/runtime/example_config/connectors/.

Example plugins

  • Structure / polling loop: postgres_source
  • External client + batching: influxdb_source
  • Arrow batch handling (sink side, for reference): iceberg_sink

Acceptance criteria

  • fluss_source plugin builds as cdylib and loads in connectors runtime
  • Log-table scan reads batches and publishes to configured Iggy stream/topic
  • Offset checkpointing survives connector restart
  • Config documented in plugin README.md with minimal runnable example
  • Unit tests for config parsing and row-to-JSON mapping; integration test or documented manual test plan for live Fluss cluster
  • cargo fmt, cargo sort --no-format, cargo clippy, cargo test -p iggy_connector_fluss_source pass

References

Alternatives considered

  1. Flink-only integration — Fluss already has a Flink connector; does not help users who standardize on Iggy connectors runtime.
  2. Generic HTTP/gRPC source — Fluss is not HTTP-first; official Rust client is the supported path.
  3. Java client via JNI — Iggy connectors are Rust cdylib plugins; fluss-rs matches the stack.

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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions