# How to Implement Robust Error Handling and Retry Mechanisms in Pathway Pipelines

> Implement robust error handling and retry mechanisms in Pathway pipelines. Learn how to use AsyncRetryStrategy and async_executor with back-off policies for automatic transient failure management.

- Repository: [Pathway/pathway](https://github.com/pathwaycom/pathway)
- Tags: how-to-guide
- Published: 2026-03-06

---

**Pathway provides a composable retry system built on `AsyncRetryStrategy` and `async_executor` that wraps UDFs, LLM calls, and async transformers with configurable back-off policies to handle transient failures automatically.**

Pathway is a Python-based data processing framework designed for both batch and streaming pipelines. Implementing robust error handling and retry mechanisms in Pathway pipelines ensures your data workflows remain resilient against transient network failures, rate limits, and temporary service outages. The framework provides a unified API through the `pathway.udfs` package that decouples retry logic from business logic.

## Core Components of Pathway's Retry System

Pathway's retry architecture rests on three composable building blocks defined in [`python/pathway/udfs.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/udfs.py) and [`python/pathway/internals/udfs/retries.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/retries.py).

### AsyncRetryStrategy Abstract Base

The **`AsyncRetryStrategy`** abstract class defines the contract for any retry policy in Pathway. Every concrete strategy must implement `await strategy.invoke(func, *args, **kwargs)`, which the async executor calls for every UDF execution. This abstraction allows the framework to treat exponential back-off, fixed delays, and no-retry policies interchangeably across the entire pipeline.

### Concrete Strategy Implementations

Pathway ships with three ready-to-use strategies in [`python/pathway/internals/udfs/retries.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/retries.py):

- **`ExponentialBackoffRetryStrategy`** – Increases wait time between attempts exponentially to avoid overwhelming struggling services. Ideal for HTTP 429 or 502 errors.
- **`FixedDelayRetryStrategy`** – Waits a constant duration between retries. Suitable for idempotent quick calls where predictable latency matters.
- **`NoRetryStrategy`** – Fails immediately on any exception. Use this for non-recoverable validation errors to avoid masking bugs.

### async_options Decorator

The **`async_options`** helper in [`python/pathway/internals/udfs/executors.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/executors.py) composes capacity limiting, timeout handling, caching, and retry strategies into a single callable. The `async_executor` uses this decorator to inject policies into the execution pipeline, ensuring every wrapped function benefits from the configured resilience mechanisms.

## How the Retry Flow Works

Understanding the execution flow helps debug retry behavior in complex pipelines. When you apply a retry strategy, Pathway orchestrates the following sequence:

1. **Function decoration** – You write a sync or async function and decorate it with `@pw.udf`.
2. **Executor injection** – The decorator receives an executor, typically `pw.udfs.async_executor(...)`.
3. **Policy wrapping** – The executor builds the final callable by wrapping the original function with `async_options`.
4. **Strategy application** – If a `retry_strategy` is supplied, `async_options` applies `with_retry_strategy`, which delegates to `AsyncRetryStrategy.invoke`.
5. **Execution loop** – The chosen strategy attempts the call, catches exceptions, sleeps for the computed back-off, and retries until the max retry count is reached.

All components are type-checked with `@check_arg_types` and fully async-compatible, ensuring they work transparently in both batch and streaming contexts.

## Where Retries Are Applied in Pathway

The retry system is not limited to simple UDFs. Pathway applies these policies consistently across the entire data-processing stack:

### User-Defined Functions (UDFs)

Any user-defined function executed on a table via `pw.this` can be retried. The retry policy wraps the underlying callable before it processes table rows.

### LLM X-Packs

Large Language Model integrations such as **`OpenAIChat`**, **`CohereChat`**, and document stores accept a `retry_strategy` argument in their constructors. This propagates the policy down to the underlying HTTP calls, automatically handling rate limits from providers like OpenAI.

### Async Transformer Connector

The **`AsyncTransformer`** class in [`python/pathway/stdlib/utils/async_transformer.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/stdlib/utils/async_transformer.py) uses the `set_options` method to forward the `retry_strategy` to `udfs.async_options`. Consequently, every row-wise async transform automatically benefits from retries without additional boilerplate.

## Best Practice Guidelines for Retry Configuration

Choosing the right strategy depends on the failure mode and operational requirements:

| Situation | Recommended Strategy |
|-----------|---------------------|
| Transient network errors (HTTP 429, 502) | `ExponentialBackoffRetryStrategy(max_retries=5, initial_delay=500, backoff_factor=2, jitter_ms=200)` |
| Idempotent quick calls where latency matters | `FixedDelayRetryStrategy(max_retries=2, delay_ms=300)` |
| Non-recoverable errors (validation failures) | `NoRetryStrategy()` – fail fast to avoid masking bugs |
| Heavy-weight ops that should not overload the service | Combine a retry strategy with capacity limiting: `pw.udfs.async_executor(capacity=4, retry_strategy=...)` |

## Code Examples

### Basic UDF with Exponential Backoff

Define a retry strategy once and reuse it across multiple UDFs to handle flaky HTTP endpoints:

```python
import pathway as pw

# Configure the retry policy

retry = pw.udfs.ExponentialBackoffRetryStrategy(
    max_retries=6,          # up to 6 attempts

    initial_delay=1_000,    # 1 s first wait

    backoff_factor=2,       # double each time

    jitter_ms=300           # random jitter to avoid thundering herd

)

@pw.udf(
    executor=pw.udfs.async_executor(
        capacity=3,                 # max 3 concurrent calls

        retry_strategy=retry        # inject the retry logic

    )
)
async def fetch_user_profile(user_id: int) -> dict:
    # Simulated external HTTP request that may fail

    response = await pw.io.http.get(f"https://api.example.com/users/{user_id}")
    response.raise_for_status()          # will raise on 5xx/4xx

    return response.json()

```

Under the hood, `async_executor` builds a callable that applies `with_timeout` (if set), then `with_retry_strategy` (the `retry` object), and finally `with_capacity`. The `ExponentialBackoffRetryStrategy.invoke` method (lines 58–99 in [`retries.py`](https://github.com/pathwaycom/pathway/blob/main/retries.py)) performs the retry loop with jittered back-off calculation.

### Adding Retries to an LLM Chat Model

LLM providers often enforce rate limits. Pass a retry strategy directly to the model constructor:

```python
from pathway.xpacks.llm import llms

# Reuse the same retry policy

chat = llms.OpenAIChat(
    model="gpt-4o-mini",
    retry_strategy=pw.udfs.ExponentialBackoffRetryStrategy(max_retries=4)
)

@pw.udf(executor=pw.udfs.async_executor())
def enrich_with_sentiment(text: str) -> str:
    # Call the LLM; failures (e.g., rate-limit) are retried automatically

    answer = chat.ask(prompt=f"Sentiment of: {text}")
    return answer

```

The `llms.OpenAIChat` constructor forwards the `retry_strategy` to the internal HTTP client, which uses `async_options` (lines 387–425 in [`executors.py`](https://github.com/pathwaycom/pathway/blob/main/executors.py)) to wire the retry logic.

### Async Transformer with Per-Row Retry

Process table rows asynchronously with automatic retry on each individual operation:

```python
import pathway as pw
from pathway.io import async_transformer

# Define a simple per-row async function

async def classify_row(text: str) -> str:
    # Could be a call to a remote classification service

    return await mock_remote_classifier(text)

# Build the async transformer and inject retry

transformer = async_transformer.AsyncTransformer(
    transformer=classify_row,
    retry_strategy=pw.udfs.FixedDelayRetryStrategy(max_retries=3, delay_ms=500)
)

# Apply it to a table

tbl = pw.debug.table_from_markdown(
    """
    id | text
    1  | "I love Pathway"
    2  | "Errors are annoying"
    """
)
result = tbl.select(classification=transformer(pw.this.text))
pw.debug.compute_and_print(result)

```

Inside `AsyncTransformer.run`, the `set_options` method (lines 86–100 in [`async_transformer.py`](https://github.com/pathwaycom/pathway/blob/main/async_transformer.py)) calls `udfs.async_options` with the supplied `retry_strategy`, ensuring every row execution benefits from the fixed-delay policy.

## Summary

- **Pathway decouples retry logic from business logic** through the `AsyncRetryStrategy` abstraction and `async_options` decorator.
- **Three concrete strategies** are available: `ExponentialBackoffRetryStrategy` for transient network errors, `FixedDelayRetryStrategy` for predictable latency, and `NoRetryStrategy` for fail-fast behavior.
- **Implementation occurs via `async_executor`** – you supply a `retry_strategy` argument when defining UDFs, and Pathway automatically wraps the callable.
- **LLM x-packs and async transformers** accept the same retry strategy objects, providing a single source of truth for resilience across HTTP calls, database queries, and model inference.
- **Source code locations**: Core definitions live in [`python/pathway/udfs.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/udfs.py), implementations in [`python/pathway/internals/udfs/retries.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/retries.py), and executor wiring in [`python/pathway/internals/udfs/executors.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/executors.py).

## Frequently Asked Questions

### What is the default retry strategy in Pathway?

By default, Pathway uses `ExponentialBackoffRetryStrategy` when you specify `retry_strategy` without explicit parameters, though many connectors default to `NoRetryStrategy` unless configured. You must explicitly pass a strategy to `async_executor` or model constructors to enable retries.

### How do I handle non-recoverable errors in Pathway UDFs?

Use **`NoRetryStrategy()`** for validation errors or business logic failures where retrying would never succeed. This strategy fails fast immediately, preventing wasted compute cycles and ensuring errors surface quickly for debugging.

### Can I combine retry strategies with rate limiting?

Yes. Pass both `capacity` and `retry_strategy` to `pw.udfs.async_executor()`. The framework applies capacity limiting (rate limiting) and retries as composable layers: `async_options` wraps your function with timeout, then retry logic, then capacity control.

### Where is the retry logic implemented in the Pathway source code?

The retry logic is implemented in [`python/pathway/internals/udfs/retries.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/retries.py) (strategy implementations), [`python/pathway/internals/udfs/executors.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/internals/udfs/executors.py) (the `async_options` decorator that wires strategies), and exposed publicly through [`python/pathway/udfs.py`](https://github.com/pathwaycom/pathway/blob/main/python/pathway/udfs.py). The `AsyncRetryStrategy.invoke` method contains the core try-catch-sleep-retry loop.