Pipeline

Pipeline

The Pipeline class provides a flexible and expressive framework for building complex data transformation and query pipelines for Firestore.

A pipeline takes data sources, such as Firestore collections or collection groups, and applies a series of stages that are chained together. Each stage takes the output from the previous stage (or the data source) and produces an output for the next stage (or as the final output of the pipeline).

Expressions can be used within each stage to filter and transform data through the stage.

NOTE: The chained stages do not prescribe exactly how Firestore will execute the pipeline. Instead, Firestore only guarantees that the result is the same as if the chained stages were executed in order.

Usage Examples:

Constructor

new Pipeline()

Example
```typescript
const db: Firestore; // Assumes a valid firestore instance.

// Example 1: Select specific fields and rename 'rating' to 'bookRating'
const results1 = await db.pipeline()
    .collection('books')
    .select('title', 'author', field('rating').as('bookRating'))
    .execute();

// Example 2: Filter documents where 'genre' is 'Science Fiction' and 'published' is after 1950
const results2 = await db.pipeline()
    .collection('books')
    .where(and(field('genre').equal('Science Fiction'), field('published').greaterThan(1950)))
    .execute();

// Example 3: Calculate the average rating of books published after 1980
const results3 = await db.pipeline()
    .collection('books')
    .where(field('published').greaterThan(1980))
    .aggregate(average(field('rating')).as('averageRating'))
    .execute();
```

Methods

delete()

Returns:
Type Description

A new {@code Pipeline} object with this stage appended to the stage list.

Example
```typescript
// Deletes all documents in the "books" collection.
firestore.pipeline().collection("books")
   .delete();
```

execute(pipelineExecuteOptions)

Executes this pipeline and returns a Promise to represent the asynchronous operation.

The returned Promise can be used to track the progress of the pipeline execution and retrieve the results (or handle any errors) asynchronously.

The pipeline results are returned in a `PipelineSnapshot` object, which contains a list of `PipelineResult` objects. Each `PipelineResult` typically represents a single key/value map that has passed through all the stages of the pipeline, however this might differ depending on the stages involved in the pipeline. For example:

  • If there are no stages or only transformation stages, each `PipelineResult` represents a single document.
  • If there is an aggregation, only a single `PipelineResult` is returned, representing the aggregated results over the entire dataset .
  • If there is an aggregation stage with grouping, each `PipelineResult` represents a distinct group and its associated aggregated values.
Parameters:
Name Type Description
pipelineExecuteOptions

Optionally specify pipeline execution behavior.

Returns:
Type Description

A Promise representing the asynchronous pipeline execution.

Example
```typescript
const futureResults = await firestore.pipeline().collection('books')
    .where(greaterThan(field('rating'), 4.5))
    .select('title', 'author', 'rating')
    .execute();
```

findNearest(options)

Performs a vector proximity search on the documents from the previous stage, returning the K-nearest documents based on the specified query vectorValue and distanceMeasure. The returned documents will be sorted in order from nearest to furthest from the query vectorValue.

Parameters:
Name Type Description
options

An object that specifies required and optional parameters for the stage.

Returns:
Type Description

A new Pipeline object with this stage appended to the stage list.

Example
```typescript
// Find the 10 most similar books based on the book description.
const bookDescription = "Lorem ipsum...";
const queryVector: number[] = ...; // compute embedding of `bookDescription`

firestore.pipeline().collection("books")
    .findNearest({
      field: 'embedding',
      vectorValue: queryVector,
      distanceMeasure: 'euclidean',
      limit: 10,                        // optional
      distanceField: 'computedDistance' // optional
    });
```

rawStage(name, params, options)

Adds a raw stage to the pipeline.

This method provides a flexible way to extend the pipeline's functionality by adding custom stages. Each raw stage is defined by a unique `name` and a set of `params` that control its behavior.

Parameters:
Name Type Description
name

The unique name of the raw stage to add.

params

A list of parameters to configure the raw stage's behavior.

options

An object of key value pairs that specifies optional parameters for the stage.

Returns:
Type Description

A new Pipeline object with this stage appended to the stage list.

Examples
Assuming there is no 'where' stage available in SDK:
```typescript
// Assume we don't have a built-in 'where' stage
firestore.pipeline().collection('books')
    .rawStage('where', [field('published').lessThan(1900)]) // Custom 'where' stage
    .select('title', 'author');
```
Parameters:
Name Type Description
options

An object that specifies required and optional parameters for the stage.

Returns:
Type Description

A new Pipeline object with this stage appended to the stage list.

Example
```typescript
db.pipeline().collection('restaurants').search({
  query: documentMatches('breakfast')
})
```

stream() → {Stream.<PipelineResult>}

Executes this pipeline and streams the results as PipelineResults.

Returns:
Type Description
Stream.<PipelineResult>

A stream of PipelineResult.

Example
```typescript
firestore.pipeline().collection('books')
    .where(greaterThan(field('rating'), 4.5))
    .select('title', 'author', 'rating')
    .stream()
    .on('data', (pipelineResult) => {})
    .on('end', () => {});
```

toArrayExpression()

Converts this Pipeline into an expression that evaluates to an array of map (objects), where each result document of the pipeline is represented as a map in the returned array.

Result Unwrapping:

  • If the items have a single field, their values are unwrapped and returned directly in the array.
  • If the items have multiple fields, they are returned as objects in the array.
Returns:
Type Description

An Expression representing the execution of this pipeline.

Example
```typescript
// Get a list of reviewers for each book
db.pipeline().collection("books")
    .define(field("id").as("current_book_id"))
    .addFields(
        db.pipeline().collection("reviews")
            .where(field("book_id").equal(variable("current_book_id")))
            .select(field("reviewer"))
            .toArrayExpression()
            .as("reviewers")
    );
```

Output:
```json
[
  {
    "id": "1",
    "title": "1984",
    "reviewers": ["Alice", "Bob"]
  }
]
```

Multiple Fields:
```typescript
// Get a list of reviews (reviewer and rating) for each book
db.pipeline().collection("books")
    .define(field("id").as("current_book_id"))
    .addFields(
        db.pipeline().collection("reviews")
            .where(field("book_id").equal(variable("current_book_id")))
            .select(field("reviewer"), field("rating"))
            .toArrayExpression()
            .as("reviews")
   );
```

Output:
```json
[
  {
    "id": "1",
    "title": "1984",
    "reviews": [
      { "reviewer": "Alice", "rating": 5 },
      { "reviewer": "Bob", "rating": 4 }
    ]
  }
]
```

toScalarExpression()

Converts this Pipeline into an expression that evaluates to a single scalar result.

Runtime Validation: The runtime validates that the result set contains zero or one item. If zero items, it evaluates to `null`.

Result Unwrapping:

  • If the item has a single field, its value is unwrapped and returned directly.
  • If the item has multiple fields, they are returned as an object.
Returns:
Type Description

An Expression representing the execution of this pipeline.

Example
```typescript
// Calculate average rating for a restaurant
db.pipeline().collection("restaurants")
    .define(field("id").as("current_restaurant_id"))
    .addFields(
      db.pipeline().collection("reviews")
        .where(field("restaurant_id").equal(variable("current_restaurant_id")))
        .aggregate(average("rating").as("avg"))
        // Unwraps the single "avg" field to a scalar double
        .toScalarExpression().as("average_rating")
   );
```

Output:
```json
{
  "name": "The Burger Joint",
  "average_rating": 4.5
}
```

Multiple Fields:
```typescript
// Calculate average rating AND count for a restaurant
db.pipeline().collection("restaurants")
    .define(field("id").as("current_restaurant_id"))
    .addFields(
      db.pipeline().collection("reviews")
        .where(field("restaurant_id").equal(variable("current_restaurant_id")))
        .aggregate(
          average("rating").as("avg"),
          count().as("count")
        )
        // Returns an object with "avg" and "count" fields
        .toScalarExpression().as("stats")
   );
```

Output:
```json
{
  "name": "The Burger Joint",
  "stats": {
    "avg": 4.5,
    "count": 100
  }
}
```