Define queues
from polos import queue
# Limit concurrent API calls
api_queue = queue("api-calls", concurrency_limit=5)
# Limit database connections
db_queue = queue("database-ops", concurrency_limit=10)
# Limit CPU-intensive work
heavy_queue = queue("heavy-processing", concurrency_limit=2)
import { Queue } from '@polos/sdk';
// Limit concurrent API calls
const apiQueue = new Queue('api-calls', { concurrencyLimit: 5 });
// Limit database connections
const dbQueue = new Queue('database-ops', { concurrencyLimit: 10 });
// Limit CPU-intensive work
const heavyQueue = new Queue('heavy-processing', { concurrencyLimit: 2 });
Assign workflows to queues
from polos import workflow, WorkflowContext
@workflow(id="api_call", queue=api_queue)
async def api_call(ctx: WorkflowContext, payload):
return await ctx.step.run("request", make_api_request, payload)
@workflow(id="db_read", queue=db_queue)
async def db_read(ctx: WorkflowContext, payload):
return await ctx.step.run("query", execute_query, payload)
# Multiple workflows can share the same queue
@workflow(id="db_write", queue=db_queue)
async def db_write(ctx: WorkflowContext, payload):
return await ctx.step.run("insert", insert_data, payload)
import { defineWorkflow } from '@polos/sdk';
const apiCallWorkflow = defineWorkflow<ApiPayload, unknown, ApiResult>(
{ id: 'api_call', queue: apiQueue },
async (ctx, payload) => {
const result = await ctx.step.run(
'make_request',
() => makeApiRequest(payload.url, payload.method ?? 'GET', payload.data),
);
return { url: payload.url, method: payload.method ?? 'GET', result };
},
);
const dbReadWorkflow = defineWorkflow<DbReadPayload, unknown, Record<string, unknown>>(
{ id: 'db_read', queue: dbQueue },
async (ctx, payload) => {
const results = await ctx.step.run(
'execute_query',
() => executeDbQuery(payload.table, payload.query ?? {}),
);
return { table: payload.table, results };
},
);
// Multiple workflows can share the same queue
const dbWriteWorkflow = defineWorkflow<DbWritePayload, unknown, Record<string, unknown>>(
{ id: 'db_write', queue: dbQueue },
async (ctx, payload) => {
const inserted = await ctx.step.run(
'insert_data',
() => insertDbData(payload.table, payload.data),
);
return { table: payload.table, inserted };
},
);
Inline queue config
@workflow(id="inline_queue", queue={"concurrency_limit": 3})
async def inline_queue(ctx: WorkflowContext, payload):
return {"message": "Processed"}
@workflow(id="named_queue", queue="my-queue")
async def named_queue(ctx: WorkflowContext, payload):
return {"message": "Processed"}
const inlineQueueWorkflow = defineWorkflow<Record<string, unknown>, unknown, Record<string, unknown>>(
{ id: 'inline_queue_workflow', queue: { name: 'inline_queue_workflow', concurrencyLimit: 3 } },
async () => {
return { message: 'Processed with inline queue' };
},
);
const namedQueueWorkflow = defineWorkflow<Record<string, unknown>, unknown, Record<string, unknown>>(
{ id: 'named_queue_workflow', queue: 'my-named-queue' },
async () => {
return { message: 'Processed with named queue' };
},
);
Run it
git clone https://github.com/polos-dev/polos.git
cd polos/python-examples/12-shared-queues
cp .env.example .env # Add your POLOS_PROJECT_ID and API key
uv sync
python main.py
git clone https://github.com/polos-dev/polos.git
cd polos/typescript-examples/12-shared-queues
cp .env.example .env # Add your POLOS_PROJECT_ID and API key
npm install
npx tsx main.ts