Key concepts#

What are datasets and blocks?#

Ray Data has two main concepts, datasets and blocks.

A Dataset represents a distributed collection of data and defines the operations that load and process it. It’s the primary API you work with in Ray Data. You typically create a Dataset from external storage or in-memory data, apply transformations to the data, and then write the outputs to external storage or feed them to training workers.

The Dataset API is lazy. Ray Data doesn’t execute operations until you materialize or consume the dataset with a method such as show(). Because execution waits until then, Ray Data can optimize the execution plan and execute operations in a pipelined, streaming fashion.

A block is a set of rows that represents a single partition of the dataset. Blocks hold their rows in columnar formats such as Arrow, and they’re the basic unit of data processing in Ray Data. Ray Data processes a dataset as follows:

  1. Ray Data partitions every dataset into a number of blocks.

  2. Ray Data distributes and parallelizes processing of the whole dataset at the block level. It processes blocks in parallel and, for the most part, independently.

The following figure shows a dataset with three blocks, each holding 1000 rows. Ray Data holds the Dataset on the process that triggers execution. That process is usually the entrypoint of the program, called the driver. Ray Data stores the blocks as objects in Ray’s shared-memory object store. Internally, Ray Data can natively handle a block as either a pandas DataFrame or a PyArrow Table.

A ray.data.Dataset holds a table that maps row ranges 1-1000, 1001-2000, and 2001-3000 to object references. Each reference points to a block of 1000 rows with the columns col1 and col2.

What are operators and plans?#

Ray Data uses a two-phase planning process to execute operations efficiently. When you write a program with the Dataset API, Ray Data first builds a logical plan, which is a high-level description of what operations to perform. When execution begins, Ray Data converts the logical plan into a physical plan that specifies exactly how to execute those operations.

The following diagram shows the complete planning process.

The LogicalOptimizer turns a logical plan into an optimized logical plan. The Planner converts the optimized logical plan into a physical plan, and the PhysicalOptimizer turns that into an optimized physical plan.

Operators are the building blocks of these plans. Ray Data uses two kinds of operators, one for each plan:

  • Logical plans consist of logical operators that describe what operation to perform. For example, when you write dataset = ray.data.read_csv(...), Ray Data creates a Read logical operator to specify what data to read.

  • Physical plans consist of physical operators that describe how to execute the operation. For example, Ray Data converts the Read logical operator into two physical operators, an InputDataBuffer and a TaskPoolMapOperator. The TaskPoolMapOperator launches Ray tasks to read the data.

The following example shows how Ray Data builds a logical plan. As you chain operations, Ray Data constructs the logical plan behind the scenes:

import ray

dataset = ray.data.range(100)
dataset = dataset.add_column("test", lambda x: x["id"] + 1)
dataset = dataset.select_columns("test")

You can inspect the resulting logical plan by printing the dataset:

Project
+- MapBatches(add_column)
   +- Dataset(schema={...})

When execution begins, Ray Data optimizes the logical plan, translates it into a physical plan, and then optimizes the physical plan. The physical plan is a series of operators that implement the data transformations. A single logical operator can become multiple physical operators during translation. For example, Read becomes both InputDataBuffer and TaskPoolMapOperator. The physical optimization pass then applies rules such as FuseOperators, which combines map operators to reduce serialization overhead.

Physical operators do the following:

  • Take in a stream of block references.

  • Perform their operation, either by transforming data with Ray tasks or actors, or by manipulating references.

  • Output another stream of block references.

For more details on Ray tasks and actors, see Ray Core Concepts.

Note

A dataset’s execution plan only runs when you materialize or consume the dataset through operations such as show().

How does streaming execution work?#

Ray Data can stream data through a pipeline of operators to process large datasets efficiently.

With streaming execution, different operators in an execution can scale independently while they run concurrently, which makes resource allocation more flexible and fine-grained. For example, if two map operators require different amounts or types of resources, the streaming execution model can run them concurrently and independently while maintaining high performance.

Note

Streaming is primarily useful for non-shuffle operations. Shuffle operations such as ds.sort() and ds.groupby() require materializing data, which stops streaming until the shuffle is complete.

The following example shows how streaming execution works in Ray Data:

import ray

def cpu_function(row):
    return row

class GPUClass:
    def __call__(self, row):
        return row

def cpu_function2(row):
    return row

def filter_func(row):
    return True

# Create a dataset with 1K rows
ds = ray.data.range(1000)

# Define a pipeline of operations
ds = ds.map(cpu_function, num_cpus=2)
ds = ds.map(GPUClass, num_gpus=1)
ds = ds.map(cpu_function2, num_cpus=4)
ds = ds.filter(filter_func)

# Data starts flowing when you call a method like show()
ds.show(5)

This code creates a logical plan like the following:

Filter(filter_func)
+- Map(cpu_function2)
   +- Map(GPUClass)
      +- Map(cpu_function)
            +- Dataset(schema={...})

The streaming topology looks like the following:

A streaming topology of Read, Map, Map, Map, and Filter operators. Each operator writes to a queue that feeds the next operator. A legend marks each queue as an out-queue for its operator and an in-queue for the next one.

In the streaming execution model, operators form a pipeline, and each operator’s output queue feeds directly into the input queue of the next downstream operator. This design creates an efficient flow of data through the execution plan.

Because of the pipeline, multiple stages can execute concurrently, which improves overall performance and resource utilization. For example, if the map operator requires GPU resources, the streaming execution model can execute the map operator concurrently with the filter operator, which might run on CPUs. This way, the pipeline uses the GPU effectively through its entire duration.

For more about the streaming execution model, see this Anyscale blog post on streaming execution across CPUs and GPUs.