What is YTsaurus Flow?

YTsaurus Flow is a framework for streaming cross-DC event processing with exactly-once guarantees within the YTsaurus ecosystem. It offers APIs for C++, Java and Kotlin, Python, and Go, and supports declarative pipeline descriptions in YQL.

Its closest external counterparts are Google Cloud Dataflow and Apache Flink.

The system is under active development, but more than ten teams have already built their production processes on it. The system can reliably:

  • Handle loads exceeding 100 GB/s or 1 million events per second.
  • Support 150+ logical pipeline nodes.

Contacts

For any issues with YTsaurus Flow, open a GitHub Issue.

System properties

  • Native support for multi-stage pipelines. As a result, you get simpler deployment and system management.
  • Support for watermarks and timers.
  • Exactly-once semantics for event processing by default.
  • Typical event processing latency under stable operation: 1s–10s.
  • Automatic balancing of partitions across machines.
  • Fault tolerance: the pipeline survives the failure of individual machines and data centers.
  • Ability to implement business logic in C++, Java and Kotlin, Python, Go, and YQL.
  • Support for stateful processing with persistent state in YTsaurus dynamic tables.
  • Support for running in YTsaurus.

Choose a language

Flow supports several languages for implementing business logic:

  • C++ — native implementation, maximum performance, full control. Use this for high-load pipelines.
  • Java and Kotlin — run via the companion mechanism. They support Spring Boot. These are suitable for teams with a JVM stack.
  • Python — runs via the companion mechanism. This is the easiest way to prototype a pipeline or process a small data stream.
  • Go — runs via the companion mechanism. A single binary runs the pipeline and acts as a companion in the job. Suitable for teams with a Go stack.
  • YQL — declarative pipeline description as an SQL query. It has a low entry barrier and doesn’t require writing code in C++, Java, Kotlin, Go, or Python. It’s under active development, and not all planned features are available yet.

Target system properties

  • Smart planning of the entire pipeline, taking into account CPU/RAM consumption of individual pipeline nodes and shared resources (common caches, databases, etc.) used by multiple nodes.
  • Ability to run pipelines on clusters with thousands of nodes or more.
  • Minimal downtime when nodes, DCs, or clusters fail, as well as during updates.

See also