Deploy Dramatiq workers on Vercel
Dramatiq is a distributed task processing library for Python. You declare functions as actors, send them messages, and workers run them in the background.
Deploy Dramatiq workers to Vercel with the Python
runtime, Vercel Queues, and
Vercel Functions. Vercel builds your Dramatiq broker as a
private, queue-triggered Vercel Function, so you don't need to run a long-lived
dramatiq worker process or a Redis or RabbitMQ broker.
Dramatiq projects on Vercel must declare their dependencies in
pyproject.toml:
[project]
name = "dramatiq-on-vercel"
version = "0.1.0"
requires-python = ">=3.12"
dependencies = [
"dramatiq>=2.2,<3",
"fastapi",
]This example uses FastAPI to enqueue messages. You can use any supported Python web framework for the producer.
Create a VercelQueueBroker and set it as Dramatiq's broker before you declare
any actors:
import dramatiq
from vercel.integrations.dramatiq import VercelQueueBroker
broker = VercelQueueBroker()
dramatiq.set_broker(broker)
@dramatiq.actor
def send_email(user_id: str) -> None:
deliver_email(user_id)Create a worker module that imports your tasks. The import is what declares each actor's queues on the broker:
from tasks import broker
__all__ = ["broker"]Add the web entrypoint and the Dramatiq worker to pyproject.toml:
[tool.vercel]
entrypoint = "main:app"
[[tool.vercel.subscribers]]
entrypoint = "worker"The subscriber entrypoint is a Python module import path, so use a dotted path
such as queues.worker for queues/worker.py, without the .py
suffix. Vercel builds main:app as the public web application and worker as a
private Vercel Function. Only Vercel Queues can invoke the worker function.
At build time, Vercel imports the module, reads every queue the broker declares,
and compiles the subscriber into a queue-triggered function. You don't need to
configure experimentalTriggers in vercel.json.
Import an actor into your web application and call send as you normally
would:
from fastapi import FastAPI
from tasks import send_email
app = FastAPI()
@app.post("/emails")
def enqueue_email(user_id: str):
message = send_email.send(user_id)
return {"messageId": message.message_id}Each call to send publishes a message to the actor's queue topic. Vercel
Queues then invokes the subscriber function, which runs the actor.
Use send_with_options to delay a message:
send_email.send_with_options(args=(user_id,), delay=60_000)Dramatiq delays are in milliseconds. Vercel Queues can delay a message by up to the message retention period. See Queues limits.
Use vercel dev to run the web application and Dramatiq worker locally:
vercel devvercel dev starts both your web application and the Dramatiq worker. You
don't need to run the dramatiq CLI in another terminal.
Deploy the project by connecting your Git repository or by using the Vercel CLI:
vc deployWhen your web function calls send, the broker publishes the message to the
topic that matches the actor's queue name. Vercel Queues invokes the private
subscriber function, which hands the delivery to an in-process Dramatiq worker,
runs the actor, and acknowledges the message after it succeeds.
Each Dramatiq queue maps to two topics: the queue itself and its Dramatiq delay
queue. The default queue becomes the default topic and the default_DDQ
topic. Vercel Queues delivers delayed messages and retries through the delay
topic, so the generated function subscribes to both.
Topics are partitioned by deployment ID, so a deployment consumes only the messages it published. See deployments and versioning.
Pass options to VercelQueueBroker to control naming, delivery, and the
underlying queue client:
| Option | Type | Default | Description |
|---|---|---|---|
consumer_group | str | dramatiq | Consumer group used for subscriptions and polling |
queue_name_prefix | str | No prefix | Prefix applied to Dramatiq queue names before topic sanitization |
retention | duration | Service default | Retention applied to published messages |
lease_duration | duration | Service default | Processing timeout for received messages |
requeue_delay_seconds | int | Zero seconds | Visibility delay used when Dramatiq requeues a message |
push_retry_delay_seconds | int | One second | Visibility delay used when a push delivery finds no free worker slot |
push_handoff_wait_seconds | float | 30 seconds | Maximum request-time wait for worker readiness and settlement |
use_message_id_as_idempotency_key | bool | False | Publish the Dramatiq message ID as the Queues idempotency key |
poll | bool | Push on Vercel | Force poll delivery when True or push delivery when False |
middleware | list | Dramatiq defaults | Middleware list passed to the base broker |
The broker also accepts token, region, base_url, deployment, timeout,
and headers, and forwards them to the queue client. See client
options.
Workers that share a topic and a consumer group compete for messages. Workers
that share a topic with different consumer groups each receive a copy of every
message. Set queue_name_prefix when other producers in the project publish to
topics with the same names as your Dramatiq queues.
Dramatiq's Retries middleware handles actor failures. When an actor raises,
the middleware republishes the message to the delay topic with exponential
backoff and acknowledges the original delivery. Configure it per actor:
@dramatiq.actor(max_retries=5, min_backoff=1_000, max_backoff=60_000)
def send_email(user_id: str) -> None:
deliver_email(user_id)Dramatiq defaults to 20 retries, a minimum backoff of 15 seconds, and a maximum
backoff of seven days. Set max_backoff so retries stay inside your message
retention window, because Vercel Queues can't deliver a message after it
expires.
If the function times out or crashes before the actor finishes, the processing lease expires and Vercel Queues delivers the message again. Queues provides at-least-once delivery, so actors should be idempotent.
Vercel Queues acts as the Dramatiq broker, not a result backend. To return
values from actors, add the Results middleware with the Vercel Runtime Cache
backend:
import dramatiq
from dramatiq.results import Results
from vercel.integrations.dramatiq import (
VercelQueueBroker,
VercelRuntimeCacheBackend,
)
broker = VercelQueueBroker()
broker.add_middleware(Results(backend=VercelRuntimeCacheBackend()))
dramatiq.set_broker(broker)
@dramatiq.actor(store_results=True)
def add(left: int, right: int) -> int:
return left + rightProducers then read results as they would with any Dramatiq result backend:
from dramatiq.composition import group
from tasks import add
result_group = group(add.message(x, x) for x in range(10)).run()
values = list(result_group.get_results(block=True, timeout=30_000))VercelRuntimeCacheBackend accepts namespace, name, and tags to control
where results are stored. Runtime Cache is
ephemeral and regional, so treat results as short-lived. Store anything you
need to keep in a database.
Blocking on results holds the producer function open for as long as the actors take, which counts against the function's maximum duration and its compute cost. For long jobs, have the actor write to a database and poll from the client instead.
With no topics filter, the generated function consumes every queue the broker
declares. Add a filter to split queues across separate functions:
[[tool.vercel.subscribers]]
entrypoint = "worker"
topics = ["emails*"]
[[tool.vercel.subscribers]]
entrypoint = "worker"
topics = ["reports*"]A trailing * matches by prefix. Prefer the prefix form so each function also
subscribes to its delay topic. A filter of topics = ["emails"] excludes
emails_DDQ, which means delayed messages and retries for that queue are never
delivered.
Dramatiq actors run inside Vercel Functions, so all Vercel Functions limitations apply, including maximum duration and bundle size.
- Message arguments: Actor arguments must be JSON-serializable. Vercel Queues supports messages up to 100 MB.
- Long-running processes: Vercel uses queue-triggered functions instead of a
persistent
dramatiqworker process. Worker control features that require persistent process state aren't available. - Queue management: Vercel Queues doesn't support purging or joining
queues, so
broker.flush(),broker.flush_all(), andbroker.join()raiseNotImplementedError. Use Dramatiq'sStubBrokerin unit tests instead.
For more about deploying Dramatiq on Vercel, see:
Was this helpful?