Rick Fanta
03/30/2026, 10:45 AMMarvin
03/30/2026, 10:45 AMMarvin
03/30/2026, 10:49 AMimport asyncio
from prefect.client.orchestration import get_client
from prefect.client.schemas.filters import FlowRunFilter
# Adjust these
WORK_POOL_NAME = "my-pool"
WORK_QUEUE_NAME = "my-consumer-queue"
# Which states to clear from the queue
STATE_TYPES_TO_CLEAR = ["SCHEDULED", "PENDING", "LATE"] # adjust as needed
BATCH_LIMIT = 50 # server-side bulk limit
async def drain_queue(pool_name: str, queue_name: str):
async with get_client() as client:
# Pause the queue to stop new work being picked up
q = await client.read_work_queue_by_name(name=queue_name, work_pool_name=pool_name)
await client.update_work_queue(q.id, is_paused=True)
while True:
filters = FlowRunFilter(
work_queue_name={"any_": [queue_name]},
state={"type": {"any_": STATE_TYPES_TO_CLEAR}},
)
# Bulk set state -> CANCELLED
resp = await <http://client._client.post|client._client.post>(
"/flow_runs/bulk_set_state",
json={
"flow_runs": filters.model_dump(mode="json"),
"state": {"type": "CANCELLED", "message": "Drained from work queue"},
"force": True,
"limit": BATCH_LIMIT,
},
)
results = resp.json().get("results", [])
if not results:
break
# Optional: resume the queue afterward
await client.update_work_queue(q.id, is_paused=False)
asyncio.run(drain_queue(WORK_POOL_NAME, WORK_QUEUE_NAME))
Option B: Delete all pending/scheduled runs for a work queue (removes them entirely)
import asyncio
from prefect.client.orchestration import get_client
from prefect.client.schemas.filters import FlowRunFilter
WORK_POOL_NAME = "my-pool"
WORK_QUEUE_NAME = "my-consumer-queue"
STATE_TYPES_TO_DELETE = ["SCHEDULED", "PENDING", "LATE"]
BATCH_LIMIT = 50
async def delete_from_queue(pool_name: str, queue_name: str):
async with get_client() as client:
q = await client.read_work_queue_by_name(name=queue_name, work_pool_name=pool_name)
await client.update_work_queue(q.id, is_paused=True)
while True:
filters = FlowRunFilter(
work_queue_name={"any_": [queue_name]},
state={"type": {"any_": STATE_TYPES_TO_DELETE}},
)
resp = await <http://client._client.post|client._client.post>(
"/flow_runs/bulk_delete",
json={
"flow_runs": filters.model_dump(mode="json"),
"limit": BATCH_LIMIT,
},
)
deleted = resp.json().get("deleted", [])
if not deleted:
break
# Optional: resume the queue afterward
await client.update_work_queue(q.id, is_paused=False)
asyncio.run(delete_from_queue(WORK_POOL_NAME, WORK_QUEUE_NAME))
CLI pointers (for quick inspection, not bulk clearing)
- Pause/resume a queue:
- prefect work-queue pause <name> -p <pool>
- prefect work-queue resume <name> -p <pool>
- Preview or read runs for a queue:
- prefect work-queue preview <name> -p <pool>
- prefect work-queue read-runs <name> -p <pool>
- Cancel/delete a single flow run:
- prefect flow-run cancel <id>
- prefect flow-run delete <id>
Docs
- Work pools and queues overview: https://docs.prefect.io/latest/learn/concepts/work-pools/Marvin
03/30/2026, 10:49 AM