YTsaurus Flow Goes Open Source

Introducing the stream processing framework for YTsaurus

We have open-sourced YTsaurus Flow under the Apache 2.0 license. It’s a framework for stateful real-time stream processing with exactly-once guarantees by default. Flow is built on YTsaurus queues, dynamic tables, and transactions. It takes care of state management, automatically adapts processing to the load, and recovers from failures, so developers can focus on business logic. The source code is available in the YTsaurus repository on GitHub.

Key features

  • Exactly-once by default. State, the metadata needed to provide the guarantees, and computation results are committed in a single YTsaurus transaction. To save resources, you can switch to at-least-once or at-most-once when needed.

  • Key state. Long-term state is stored in dynamic tables, and shuffle by key brings related events together.

  • Auto-balancing. Flow picks the number of partitions and distributes load across workers on its own.

  • Multi-datacenter. Pipelines keep running even if an entire data center fails or is down for planned maintenance.

  • Event time. Message timestamps, watermarks, and timers are supported natively.

  • Complex graphs. Hundreds of nodes, cycles, and topology visualization in the YTsaurus web UI.

  • Language settings. The core is written in C++, and you can write business logic in Python, Go, Java, Kotlin, or C++. Simple pipelines can be defined in YQL.

To get started, see the YTsaurus documentation and then the Flow documentation. If something is missing, let us know in the community chat or open an issue or PR in the repository.

Sign in to save this post