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

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 and 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:

  • 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 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 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:

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) 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:

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) to wire the retry logic.

Async Transformer with Per-Row Retry

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

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) 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, implementations in python/pathway/internals/udfs/retries.py, and executor wiring in 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 (strategy implementations), python/pathway/internals/udfs/executors.py (the async_options decorator that wires strategies), and exposed publicly through python/pathway/udfs.py. The AsyncRetryStrategy.invoke method contains the core try-catch-sleep-retry loop.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →