Open-source project
arkflow-rs/arkflow avatar
arkflow-rs/arkflow

ArkFlow: a Rust stream processing engine with SQL pipelines and optional AI models

High performance Rust stream processing engine seamlessly integrates AI capabilities, providing powerful real-time data processing and intelligent analysis.

1,304 stars46 forksRustApache-2.0

At a glance

What is it?
ArkFlow is an Apache-2.0 Rust engine that runs YAML-defined streams over Kafka, MQTT, SQL and file sources, processes them with DataFusion SQL, and can attach machine learning inference. This review covers the mechanism, the install path, and where the design gets in your way.
Who is it for?
Adopt ArkFlow if your pipeline is naturally expressible as a YAML stream with a SQL step and you want the whole thing in one Rust binary rather than a cluster of JVM services. Skip it if you need a documented Python UDF contract, a stable plugin ABI, or a stored-state upgrade path, because the README documents none of those.
Can I use it commercially?
Yes. Apache-2.0 is a permissive licence: you can use, modify and sell software built on it, as long as you keep its copyright and licence notices.
Is it still maintained?
Yes. The repository received new commits within the last day.
What is it written in?
Mainly Rust, according to GitHub's language statistics.

Answers come from the project's GitHub data, last synced on October 1, 2026, and from our analysis. They are not legal advice.

Editorial analysis

What ArkFlow solves, and who is actually in the target audience

Most teams that need to move events from Kafka into a database, filter them, and occasionally score them with a model end up assembling three or four services. ArkFlow's premise is that this is one program. The README describes it as a "High performance Rust stream processing engine" that "supports multiple input/output sources and processors" and can load and execute machine learning models over streaming data. The unit of work is a stream: an input, a pipeline, an output, an error_output, and optionally a buffer, all declared in a single YAML file.

The audience is narrower than the feature list suggests. If you already run a Flink or Kafka Streams topology and your team writes Java or Scala, ArkFlow is not replacing that. The people who benefit are Rust shops, small platform teams that do not want a JVM in the data path, and anyone whose transformation is genuinely a SQL SELECT over a stream. The repository layout reinforces this: a Cargo workspace with crates/arkflow-core, crates/arkflow-plugin, crates/arkflow-server and crates/arkflow, plus a console/ directory and a docker/ directory. It is a product with a CLI and a control plane, not a library you embed in three lines.

How a stream actually moves: inputs, Arrow batches, DataFusion, sinks

The mechanism is visible in the Quick Start configuration. An input produces records. The pipeline declares thread_num and a list of processors. The first processor in the README's example is json_to_arrow, which converts JSON records into Arrow columnar batches; the second is a sql processor whose query is "SELECT * FROM flow WHERE value >= 10". That ordering matters. The SQL processor is DataFusion, pinned at version 54.1 in the workspace dependencies with the avro feature enabled, and DataFusion operates on Arrow data. The table name flow is the convention the query is written against. Output then consumes whatever the pipeline emits, and error_output receives messages that failed processing.

Because the pipeline is batch-oriented rather than record-oriented, the input's batch_size and interval settings control the rhythm. The generate input in the example emits a JSON context every 1s in batches of 10. Processors that make sense over a batch include the Python UDF processor, which the README says runs user-defined functions "over the batch", and the VRL processor, which links to Vector's VRL reference. Codecs sit alongside this: JSON and Protobuf, plus Debezium CDC envelopes and Confluent Schema Registry wire format, with examples/cdc_debezium.yaml and examples/howto_cdc_schema_registry.yaml in the tree.

The durability story is the part worth reading carefully. The README states at-least-once delivery "by default via per-stream WAL durability, with optional exactly-once for transactional sinks". That is a per-stream write-ahead log, and exactly-once is conditional on the sink being transactional. examples/eos-kafka.yaml exists for that case, and there is a cluster of durability examples including S3-backed and compressed variants. If your sink is not transactional, the exactly-once option does not apply to you, and the README does not pretend otherwise.

Installing ArkFlow from source and running the generate-to-stdout example

The README documents one installation path: building from source. There is no published binary, container image or package manager command in the installation section, although a docker/ directory exists in the repository. The build clones the repository, compiles with Cargo, and runs the test suite.

bash
git clone https://github.com/arkflow-rs/arkflow.git
cd arkflow
cargo build --release
cargo test

The workspace declares rust-version = "1.97", so a toolchain at least that new is required before the build will start. The release build produces ./target/release/arkflow, which is the binary the Quick Start invokes.

The first real use is the README's own config.yaml. It generates synthetic records, converts them to Arrow, filters with SQL, and prints to standard output.

yaml
logging:
  level: info
streams:
  - input:
      type: "generate"
      context: '{ "timestamp": 1625000000000, "value": 10, "sensor": "temp_1" }'
      interval: 1s
      batch_size: 10

    pipeline:
      thread_num: 4
      processors:
        - type: "json_to_arrow"
        - type: "sql"
          query: "SELECT * FROM flow WHERE value >= 10"

    output:
      type: "stdout"
    error_output:
      type: "stdout"

Run it with the config flag, and the reader should see the generated records printed to the console once per interval, filtered by the value >= 10 predicate.

bash
./target/release/arkflow --config config.yaml

Swapping the output for Kafka means replacing the output block with type: kafka, a brokers list, a topic block, and a client_id, as the README's output example shows. The topic field is not a plain string in that example; it is a nested structure with type and value keys, which is worth noticing before you copy a topic name in as a scalar.

The parts of ArkFlow that will slow you down

The plugin story is the biggest gap. The README says the design is "Modular" and "easy to extend with new input, buffer, output, and processor components", and the workspace has a crates/arkflow-plugin member. What the README does not document is the trait contract, the registration mechanism, or whether third-party plugins can be loaded without recompiling the binary. If you are evaluating ArkFlow specifically to write a custom connector, the documentation you have does not answer that question, and you should read the source in crates/arkflow-plugin before assuming it is a plugin system in the loadable sense.

The Python UDF processor has the same shape of problem. It is listed as a processor and described as running Python functions over the batch. The README does not show a single example of how a Python function is declared, where the file lives, or how the runtime is embedded. That is a gap you would have to close from the repository or the docs site before designing around it.

Durability is a trade-off rather than a free feature. Per-stream WAL means disk writes on the path, and the examples directory contains both an aggressive variant and a compressed variant, which suggests the default is not tuned for every workload. The exactly-once option narrows your sink choices to transactional ones. And the buffer components, memory buffer and session window, sit in front of that: a memory buffer is exactly what it sounds like, and the README's own description ties it to high-throughput scenarios and window aggregation, which is also where you would lose the most on restart.

The control plane is optional and separate. The README describes a Hub and web console to "observe, configure, and operate multiple ArkFlow compute nodes as a fleet", with examples/control_plane_hub.yaml and examples/control_plane_node.yaml in the tree. If you run one node, you are carrying a feature you do not use.

ArkFlow versus Benthos and Vector for the same job

The closest comparison is Benthos, now distributed as Redpanda Connect, and Vector. Both are also single-binary stream processors with YAML configuration and a long list of connectors, and Vector's VRL is reused by ArkFlow as a processor, which tells you the authors expect overlap with that audience.

The difference is what sits in the middle. Benthos and Vector are mapping engines: you describe per-message transformations, and their blobs and mapping languages are the core abstraction. ArkFlow inserts a columnar batch and a SQL engine between input and output. The json_to_arrow processor followed by a DataFusion sql processor means you write joins, filters and aggregations as SQL against a table named flow rather than as per-record expressions. If your transformation is a join between a stream and a reference table, or a windowed aggregate, that is a real advantage. If your transformation is renaming three fields and adding a timestamp, the Arrow conversion and SQL planning are overhead you are paying for nothing.

The second difference is the AI angle. ArkFlow's README positions model loading and inference as a first-class capability, which neither Benthos nor Vector claims. The README does not document which model formats or runtimes are supported, so treat that as a direction rather than a specification you can plan against today.

Licence, maintenance and what an upgrade costs

ArkFlow is Apache-2.0, stated in the README badge, the Cargo.toml workspace package section, and the LICENSE file. Apache-2.0 is permissive with an explicit patent grant, and it imposes no copyleft on your own code. If you redistribute a modified ArkFlow, the licence requires you to carry the licence and state changes. That is a summary of the licence text, not legal advice; have counsel read it if you are shipping a derivative.

The repository is not archived, and the last push was on 2026-09-10, ten days before this writing, so calling it actively developed is defensible on the evidence. The release cadence is slower than the commit cadence: v0.5.0 landed on 2025-10-19, after v0.4.0-rc1 on 2025-06-13 and v0.3.1 on 2025-05-11. Note that the workspace version in Cargo.toml is 0.5.0, matching the newest release. There is no 1.0, and the project has already shipped a release candidate in the 0.4 line, which is normal for pre-1.0 software and also a warning that configuration keys can move.

Upgrade cost is dominated by the dependency pinning. DataFusion is pinned at 54.1, arrow-json and arrow-pyarrow at 58.4, and sqlx at 0.8.6. DataFusion's SQL surface and Arrow's format details change between minor versions, so a bump in either is a real migration, not a patch. If you vendor a fork, you own that. If you track upstream, budget for reading the release notes of each release rather than assuming the YAML config is stable. The README does not document a config migration or deprecation policy, which is the specific thing to ask about before you write forty stream definitions.

Editorial conclusion

Adopt ArkFlow if your pipeline is naturally expressible as a YAML stream with a SQL step and you want the whole thing in one Rust binary rather than a cluster of JVM services. Skip it if you need a documented Python UDF contract, a stable plugin ABI, or a stored-state upgrade path, because the README documents none of those. Before committing, build the workspace at the pinned rust-version of 1.97, run the generate-to-stdout example from the Quick Start, and read the durability examples under examples/ to see which delivery guarantee your sink actually supports.

Frequently asked questions

How do I install ArkFlow?

The README documents building from source: clone https://github.com/arkflow-rs/arkflow.git, then run cargo build --release and cargo test. The workspace declares rust-version = "1.97", so a toolchain at least that new is required. The installation section does not list a prebuilt binary or package manager command.

What data sources and outputs does ArkFlow support?

Inputs include Kafka, MQTT, HTTP, files (CSV, JSON, Parquet, Avro, Arrow), Generate, SQL databases (MySQL, PostgreSQL, SQLite), NATS, Pulsar, Redis, WebSocket, Modbus and Memory. Outputs include Kafka, MQTT, HTTP, InfluxDB 2.x, NATS, Pulsar, Redis, SQL databases, standard output and a drop sink. The error_output accepts any of the output components.

Does ArkFlow guarantee exactly-once delivery?

The README states that delivery is at-least-once by default via per-stream WAL durability, with optional exactly-once for transactional sinks. The exactly-once path therefore depends on the sink being transactional, and examples/eos-kafka.yaml is the example for that case.

What SQL engine does the ArkFlow sql processor use?

The workspace Cargo.toml pins datafusion at version 54.1 with the avro feature, and the Quick Start pipeline runs a query of the form SELECT * FROM flow WHERE value >= 10 against a table named flow. A json_to_arrow processor runs before the SQL step in that example.

How do I write a Python UDF in ArkFlow?

The README lists a Python processor that runs Python user-defined functions over the batch, but it does not show how a function is declared, where the file lives, or how the Python runtime is embedded. That detail has to come from the repository or the docs site.

Official sources

  1. arkflow-rs/arkflow on GitHub
  2. License: Apache-2.0
  3. Project website
  4. README
  5. Releases
Add this badge to your README

If you maintain this project, the badge below links readers to this analysis and shows its maintenance status from the daily GitHub snapshot. Paste the markdown into your README; add ?metric=license or ?metric=stars to the image URL for a different field.

Add this badge to your README

markdown
[![Hysen Labs](https://hysenlabs.com/badge/arkflow-rs-arkflow.svg)](https://hysenlabs.com/projects/arkflow-rs-arkflow)