Metadata-Version: 2.4
Name: flask-async-celery
Version: 0.1.2
Summary: AsyncIO execution pool for Celery with Flask integration
Author: Mazhar Ali
License: GPL-3.0-only
Keywords: celery,asyncio,flask,async,tasks,redis
Classifier: Development Status :: 3 - Alpha
Classifier: Framework :: Flask
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: GNU General Public License v3 (GPLv3)
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries
Classifier: Topic :: System :: Distributed Computing
Requires-Python: >=3.10
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: celery<5.7,>=5.6
Requires-Dist: Flask>=2.3
Requires-Dist: psutil>=5.9
Provides-Extra: redis
Requires-Dist: redis>=5.0; extra == "redis"
Provides-Extra: test
Requires-Dist: pytest>=8.0; extra == "test"
Requires-Dist: pytest-asyncio>=0.24; extra == "test"
Provides-Extra: dev
Requires-Dist: pytest>=8.0; extra == "dev"
Requires-Dist: pytest-asyncio>=0.24; extra == "dev"
Requires-Dist: build>=1.2; extra == "dev"
Requires-Dist: twine>=5.0; extra == "dev"
Dynamic: license-file

# Flask Async Celery

Run native `async def` Celery tasks on a persistent asyncio event loop with bounded concurrency, Flask integration, and Celery consumer-side backpressure.

## Features

- Persistent `asyncio` event loop in a dedicated thread per Celery worker process.
- Run native `async def` Celery tasks.
- Bounded asynchronous concurrency with `max_tasks`.
- Celery task request context propagation into the asyncio execution thread.
- `self.retry()` support for async tasks.
- Normal synchronous Celery tasks continue to work.
- Redis consumer-side backpressure through Celery's `worker_disable_prefetch`.
- Flask extension with simple configuration.
- Graceful asyncio executor shutdown.
- Compatible with Celery 5.6.x and Python 3.10+.

## Architecture

```text
                         Redis
                           │
                           ▼
                    Celery Consumer
                           │
                           │ worker_disable_prefetch
                           ▼
                      AsyncIOPool
                     max_tasks = N
                           │
                           ▼
                     Bridge Threads
                           │
                           ▼
                     AsyncExecutor
                    asyncio.Semaphore(N)
                           │
                           ▼
                  Persistent asyncio loop
                           │
               ┌───────────┼───────────┐
               ▼           ▼           ▼
           Async task  Async task  Async task
```

The package separates Celery's worker execution model from asyncio execution:

1. Celery receives and traces the task.
2. `AsyncIOPool` bridges Celery execution into the asyncio executor.
3. `AsyncExecutor` owns a persistent asyncio event loop.
4. An asyncio semaphore limits the number of actively executing async tasks.
5. Celery's Redis `worker_disable_prefetch` option can reduce unnecessary task reservation when consumer-side backpressure is enabled.

The package does **not** manually pause and resume the Celery consumer. Consumer-side backpressure is provided through Celery's supported `worker_disable_prefetch` behavior.

## Requirements

- Python 3.10+
- Celery 5.6.x
- Flask 2.3+
- Redis when using the Redis broker/result backend and consumer-side backpressure

## Installation

### From PyPI

```bash
pip install flask-async-celery
```

### With Redis support

```bash
pip install "flask-async-celery[redis]"
```

### For development and testing

```bash
pip install "flask-async-celery[test]"
```

### Development installation

Clone the repository and install it in editable mode:

```bash
git clone https://github.com/pyfuncode/flask_celery_async.git
cd flask_celery_async

pip install -e ".[test]"
```

## Basic Flask Setup

```python
import asyncio

from flask import Flask

from flask_async_celery import AsyncCelery


app = Flask(__name__)

celery = AsyncCelery(
    app,
    broker_url="redis://127.0.0.1:6379/0",
    result_backend="redis://127.0.0.1:6379/1",
    max_tasks=5,
    disable_prefetch=True,
)


@celery.task
async def my_task(value):
    await asyncio.sleep(1)
    return value * 2
```

Send the task normally:

```python
result = my_task.delay(10)

print(result.get(timeout=30))
# 20
```

`AsyncCelery.task()` automatically uses `AsyncTask` for native async task functions. You can also specify `base=AsyncTask` explicitly when you want to make the task base visible:

```python
from flask_async_celery import AsyncTask


@celery.task(base=AsyncTask)
async def another_task():
    await asyncio.sleep(1)
    return "done"
```

## Flask Configuration

Configuration can be supplied through Flask:

```python
app.config["ASYNC_CELERY_MAX_TASKS"] = 10
app.config["ASYNC_CELERY_DISABLE_PREFETCH"] = True

celery = AsyncCelery(
    app,
    broker_url="redis://127.0.0.1:6379/0",
    result_backend="redis://127.0.0.1:6379/1",
)
```

Available settings:

| Setting | Default | Description |
| --- | ---: | --- |
| `ASYNC_CELERY_MAX_TASKS` | `20` | Maximum number of concurrently executing async tasks. |
| `ASYNC_CELERY_DISABLE_PREFETCH` | `True` | Enables Celery consumer-side backpressure where supported. |

Constructor arguments can also be used directly:

```python
celery = AsyncCelery(
    app,
    max_tasks=10,
    disable_prefetch=True,
)
```

Flask configuration takes precedence over the constructor defaults when the extension is initialized.

## Worker Configuration

Run the worker using the package's custom pool:

```bash
celery -A your_app.celery worker \
    -P flask_async_celery.pool:AsyncIOPool \
    -c 5 \
    --loglevel=INFO
```

For example:

```bash
celery -A your_app.celery worker \
    -P flask_async_celery.pool:AsyncIOPool \
    -c 10 \
    --loglevel=INFO
```

The extension configures Celery's worker concurrency from `max_tasks`.

If you provide `-c` on the worker command line, make sure it matches the configured `max_tasks` value.

For example:

```python
celery = AsyncCelery(
    app,
    max_tasks=5,
)
```

should normally be started with:

```bash
celery -A your_app.celery worker \
    -P flask_async_celery.pool:AsyncIOPool \
    -c 5 \
    --loglevel=INFO
```

The custom pool exposes its configured concurrency through `num_processes`, allowing Celery's consumer to use the same capacity when consumer-side prefetch is disabled.

## Concurrency

Set the maximum number of simultaneously executing async tasks with:

```python
celery = AsyncCelery(
    app,
    max_tasks=5,
)
```

With:

```text
max_tasks = 5
```

the asyncio executor allows at most five active coroutines at once.

Additional work waits for an available execution slot.

This is different from simply creating more threads. The package uses one persistent asyncio event loop and runs async coroutines concurrently on that loop.

Each Celery worker process has its own asyncio executor and event loop.

## Consumer Backpressure

For Redis, Celery 5.6 supports:

```python
worker_disable_prefetch = True
```

The extension enables this behavior by default:

```python
celery = AsyncCelery(
    app,
    max_tasks=5,
    disable_prefetch=True,
)
```

When supported by the broker transport, this provides consumer-side backpressure in addition to the executor's own concurrency limit:

```text
Celery Consumer
      │
      │ worker_disable_prefetch
      │
      ▼
  AsyncIOPool
      │
      │ max_tasks
      ▼
 AsyncExecutor
      │
      │ semaphore
      ▼
 asyncio tasks
```

The executor's semaphore remains the final execution boundary.

You can disable the consumer-side behavior:

```python
celery = AsyncCelery(
    app,
    max_tasks=5,
    disable_prefetch=False,
)
```

When disabled, the asyncio executor still enforces its own concurrency limit.

### Redis Requirement

`worker_disable_prefetch` is intended for supported Redis worker configurations.

If you use another broker, verify that your Celery version and broker transport support this feature before relying on consumer-side backpressure.

The executor-level `max_tasks` limit remains independent of consumer prefetch behavior.

The package does not manually pause or resume the Celery consumer.

## Async Tasks

Native async tasks can be declared directly with `@celery.task`:

```python
@celery.task
async def fetch_data():
    await some_async_operation()
    return "done"
```

The extension automatically uses `AsyncTask` for async task functions.

You can also explicitly specify the task base:

```python
from flask_async_celery import AsyncTask


@celery.task(base=AsyncTask)
async def fetch_data():
    await some_async_operation()
    return "done"
```

### Celery Task Features

Async tasks can use normal Celery task features such as bound tasks and retries:

```python
@celery.task(
    bind=True,
    max_retries=3,
)
async def process_item(self, item_id):
    try:
        return await process(item_id)
    except TemporaryError as exc:
        raise self.retry(
            exc=exc,
            countdown=5,
        )
```

## Celery Request Context

The package propagates the Celery task request from the Celery worker execution thread into the asyncio execution thread.

This preserves Celery request information such as:

```python
self.request.id
self.request.retries
self.request.delivery_info
```

and allows features such as:

```python
self.retry()
```

to continue working for async tasks.

The request is pushed before async execution and removed afterward.

### Flask HTTP Request Context

Celery task request propagation is different from Flask HTTP request-context propagation.

The package does **not** keep a Flask HTTP request context alive while a background Celery task executes.

If a background task needs information from an HTTP request, pass that information explicitly as task arguments.

For example:

```python
@celery.task
async def process_user(user_id, request_id):
    ...
```

This is preferable to depending on the lifetime of the original HTTP request.

## Synchronous Tasks

Normal synchronous Celery tasks can still be used:

```python
@celery.task
def sync_task(value):
    return value * 2
```

The package does not require every task to be asynchronous.

Synchronous tasks continue through the normal Celery task execution path.

## Retries

Async `self.retry()` is supported:

```python
@celery.task(
    bind=True,
    max_retries=3,
)
async def retrying_task(self):
    if should_retry():
        raise self.retry(countdown=5)

    return "success"
```

The Celery task request is preserved when execution moves from the Celery worker thread to the asyncio event loop.

This allows Celery retry metadata and delivery information to remain available to the async task.

## Exceptions

Exceptions raised by an async task propagate through the normal Celery execution path:

```python
@celery.task
async def failing_task():
    raise RuntimeError("something went wrong")
```

Celery remains responsible for:

- task failure state
- result handling
- retry behavior
- worker-level task tracing

The package provides the asyncio execution layer without replacing Celery's task tracing and lifecycle handling.

## Graceful Shutdown

The asyncio executor runs in a dedicated daemon thread.

During normal pool shutdown, the executor:

1. Stops accepting new work.
2. Stops the asyncio event loop.
3. Cancels pending asyncio tasks.
4. Waits for the loop thread when requested.
5. Closes the asyncio event loop.

The pool and executor remain responsible for their own lifecycle cleanup.

## Public API

The main public API is intentionally small:

```python
from flask_async_celery import AsyncCelery, AsyncTask
```

### `AsyncCelery`

Provides Flask integration and Celery configuration:

```python
AsyncCelery(
    app=None,
    *,
    celery=None,
    broker_url=None,
    result_backend=None,
    max_tasks=20,
    disable_prefetch=True,
)
```

### `AsyncTask`

Base class for asynchronous Celery tasks:

```python
@celery.task(base=AsyncTask)
async def my_task():
    ...
```

In most cases you can simply use:

```python
@celery.task
async def my_task():
    ...
```

because `AsyncCelery` automatically uses `AsyncTask` for async tasks.

## Development

Clone the repository and install the project in editable mode:

```bash
git clone https://github.com/pyfuncode/flask_celery_async.git
cd flask_celery_async

pip install -e ".[test]"
```

Run the test suite:

```bash
pytest -v
```

The test suite covers:

- asyncio executor concurrency
- asyncio executor lifecycle and shutdown
- `AsyncIOPool` execution
- async task exceptions
- async retries
- synchronous task execution
- synchronous retries
- Celery worker integration
- Redis consumer backpressure
- Celery Hub wakeup handling
- Flask extension configuration
- Flask application-context isolation
- end-to-end Flask/Celery/async execution

## Project Structure

```text
flask-async-celery/
├── pyproject.toml
├── README.md
├── LICENSE
├── src/
│   └── flask_async_celery/
│       ├── __init__.py
│       ├── extension.py
│       ├── executor.py
│       ├── bootstep.py
│       ├── hub_bootstep.py
│       ├── hub_wakeup.py
│       ├── pool.py
│       └── task.py
└── test/
    ├── conftest.py
    ├── test_executor.py
    ├── test_tasks.py
    ├── test_backpressure.py
    ├── test_celery_pool.py
    ├── test_worker_integration.py
    ├── test_extension.py
    └── test_extension_integration.py
```

## Design Notes

This package does not replace Celery's task tracing and lifecycle handling.

Celery remains responsible for:

- task delivery
- task acknowledgment
- retries
- result state
- task IDs
- worker lifecycle
- task tracing

The package provides the asyncio execution layer and integrates it with Celery's pool interface.

### Persistent Event Loop

The asyncio event loop is persistent for the lifetime of the worker process rather than creating a new event loop for every task.

Creating a new event loop for every task adds unnecessary setup and teardown overhead.

Instead, each worker process owns one persistent asyncio event loop:

```text
Celery Worker Process
│
├── Celery Consumer
├── Celery task execution
├── bridge threads
│
└── asyncio event loop thread
    ├── Task A
    ├── Task B
    └── Task C
```

Async tasks can therefore share the same event loop while still being bounded by `max_tasks`.

### Request Propagation

Celery's task request is associated with the worker execution context.

The asyncio event loop runs in a separate thread, so the package explicitly transfers the current Celery request into that execution context.

This preserves information required by features such as:

```python
self.request.id
self.request.retries
self.request.delivery_info
self.retry()
```

The request is pushed before async execution and removed afterward.

This is Celery task-request propagation, not Flask HTTP request-context propagation.

## Execution Model

The execution flow can be summarized as:

```text
Celery Consumer
      │
      ▼
 AsyncIOPool
      │
      ▼
 Bridge Thread
      │
      ▼
 Celery Task
      │
      ▼
 AsyncExecutor
      │
      ▼
 Persistent asyncio Event Loop
      │
      ├── Coroutine A
      ├── Coroutine B
      └── Coroutine C
```

The bridge thread allows Celery's synchronous execution and tracing model to interact with the asynchronous execution model.

The asyncio executor then schedules the coroutine on the persistent event loop.

The bridge thread waits for the asyncio execution to complete so that Celery can continue using its normal synchronous task execution and callback model.

## Limitations

### Broker-Specific Backpressure

Consumer-side `worker_disable_prefetch` support depends on Celery and the broker transport.

The asyncio executor's own `max_tasks` limit remains the final execution boundary.

If consumer-side prefetch control is unavailable for a broker, the executor still prevents more than `max_tasks` async tasks from actively executing.

### Worker Pool

The worker must use:

```text
flask_async_celery.pool:AsyncIOPool
```

for the package's asyncio execution model.

For example:

```bash
celery -A your_app.celery worker \
    -P flask_async_celery.pool:AsyncIOPool \
    -c 5 \
    --loglevel=INFO
```

### One Event Loop Per Worker Process

Each worker process owns its own asyncio event loop and concurrency limit.

For example:

```bash
celery -A your_app.celery worker \
    -P flask_async_celery.pool:AsyncIOPool \
    -c 5 \
    --loglevel=INFO
```

creates a worker configuration with five execution slots.

If you run multiple worker processes, each process has its own pool, bridge threads, and asyncio event loop.

The total async capacity is therefore distributed across the worker processes.

### HTTP Request Data

A Celery task should not depend on the lifetime of the Flask HTTP request that originally triggered it.

Pass required request-specific information explicitly to the task.

### Task Termination

The package relies on Celery's normal worker lifecycle and task tracing mechanisms. It does not provide a separate public API for individually terminating asyncio tasks.

## Testing

Run the complete test suite:

```bash
pytest -v
```

The project tests the execution model with real Celery workers in addition to unit-level executor and pool tests.

The test suite covers:

- executor concurrency
- executor lifecycle
- async task execution
- synchronous task execution
- async exceptions
- async retries
- synchronous retries
- Celery worker integration
- Redis consumer backpressure
- Celery Hub wakeup behavior
- Flask extension configuration
- Flask application-context handling
- end-to-end Flask/Celery/async execution

A successful test run should show all tests passing.

## Version

Current version:

```text
0.1.2
```

## License

This project is licensed under the **GNU General Public License v3.0**.

Copyright (c) 2026 Mazhar Ali

This software is distributed under the terms of the GNU General Public License version 3.0.

See the [LICENSE](LICENSE) file for the complete license text.

For the full license terms, see the official GNU General Public License v3.0 text.
