Rate Limiting
Stem supports per-task rate limits via TaskOptions.rateLimit and a pluggable
RateLimiter interface. This lets you throttle hot handlers with a shared
Redis-backed limiter or custom driver.
Stem also supports group-scoped rate limits with TaskOptions.groupRateLimit
for shared quotas across multiple task types/tenants.
Quick start
- Task Options
- Worker Wiring
- Producer Enqueue
final tasks = <TaskHandler<Object?>>[
FunctionTaskHandler<void>(
name: _taskName,
options: const TaskOptions(
queue: 'throttled',
maxRetries: 0,
visibilityTimeout: Duration(seconds: 60),
rateLimit: const RateLimit.perSecond(3),
),
entrypoint: _renderEntrypoint,
),
];
final worker = Worker(
broker: broker,
tasks: tasks,
backend: backend,
rateLimiter: rateLimiter,
queue: 'throttled',
consumerName:
Platform.environment['WORKER_NAME'] ?? 'rate-limit-demo-worker',
subscription: RoutingSubscription.singleQueue('throttled'),
concurrency: 2,
);
for (var i = 0; i < totalJobs; i++) {
final delaySeconds = i >= totalJobs / 2 ? 4 : 0;
final notBefore = delaySeconds > 0
? DateTime.now().add(Duration(seconds: delaySeconds))
: null;
final priority = i.isEven ? 9 : 2;
final route = routing.resolve(
RouteRequest(
task: taskName(),
headers: const {},
queue: 'throttled',
),
);
final appliedPriority = route.effectivePriority(priority);
final id = await client.enqueue(
taskName(),
args: {
'job': i + 1,
'scheduledFor': notBefore?.toIso8601String(),
'requestedPriority': priority,
},
options: TaskOptions(
queue: 'throttled',
priority: priority,
maxRetries: 0,
),
notBefore: notBefore,
meta: {
'requestedPriority': priority,
'appliedPriority': appliedPriority,
if (notBefore != null) 'scheduledFor': notBefore.toIso8601String(),
},
);
stdout.writeln(
'[producer] job=${i + 1} priority=$priority '
'applied=$appliedPriority delay=${delaySeconds}s id=$id',
);
}
Docs snippet (in-memory demo)
- Define a rate-limited task
- Limiter config + state
- Limiter acquire decision
- Wire worker with rate limiter
- Enqueue with tenant header
- Bootstrap StemApp
- Start worker
- Create Stem client
- Enqueue demo task
- Shutdown cleanly
class RateLimitedTask extends TaskHandler<void> {
String get name => 'demo.rateLimited';
TaskOptions get options => const TaskOptions(
rateLimit: const RateLimit.perSecond(10),
maxRetries: 3,
);
Future<void> call(TaskContext context, Map<String, Object?> args) async {
final actor = args['actor'] as String? ?? 'anonymous';
print('Handled rate-limited task for $actor');
}
}
DemoRateLimiter({required this.capacity, required this.interval});
final int capacity;
final Duration interval;
int _used = 0;
DateTime _windowStart = DateTime.now();
Future<RateLimitDecision> acquire(
String key, {
int tokens = 1,
Duration? interval,
Map<String, Object?>? meta,
}) async {
final window = interval ?? this.interval;
final now = DateTime.now();
final elapsed = now.difference(_windowStart);
if (elapsed >= window) {
_windowStart = now;
_used = 0;
}
if (_used + tokens <= capacity) {
_used += tokens;
return RateLimitDecision(allowed: true, meta: {'key': key});
}
final retryAfter = window - elapsed;
return RateLimitDecision(
allowed: false,
retryAfter: retryAfter.isNegative ? Duration.zero : retryAfter,
meta: {'key': key},
);
}
final limiter = DemoRateLimiter(
capacity: 2,
interval: const Duration(seconds: 1),
);
final workerConfig = StemWorkerConfig(rateLimiter: limiter);
Future<String> enqueueRateLimited(TaskEnqueuer stem) async {
return stem.enqueue(
'demo.rateLimited',
args: {'actor': 'acme'},
headers: const {'tenant': 'acme'},
);
}
// #region rate-limit-worker
final limiter = DemoRateLimiter(
capacity: 2,
interval: const Duration(seconds: 1),
);
final workerConfig = StemWorkerConfig(rateLimiter: limiter);
// #endregion rate-limit-worker
final app = await StemApp.inMemory(
tasks: [RateLimitedTask()],
workerConfig: workerConfig,
);
// #region rate-limit-demo-worker-start
await app.start();
// #endregion rate-limit-demo-worker-start
await app.start();
final stem = app;
await enqueueRateLimited(stem);
await app.close();
Run the rate_limit_delay example for a full demo:
packages/stem/example/rate_limit_delay
Rate limit values
In Dart code, use the typed RateLimit value object:
const TaskOptions(
rateLimit: RateLimit.perMinute(100),
groupRateLimit: RateLimit.perSecond(5),
)
String values remain supported at JSON/YAML and environment-configuration boundaries:
10/s— 10 tokens per second100/m— 100 tokens per minute500/h— 500 tokens per hour
groupRateLimit uses the same syntax. The worker receives a validated
RateLimit value rather than parsing strings during task execution.
How it works
- The worker asks the configured limiter to acquire the typed
rateLimit. - The worker asks the
RateLimiterfor an acquire decision. - If denied, the task is retried with backoff and
rateLimited=truemetadata. - Retry delays come from the limiter
retryAfterif provided, otherwise the worker’s retry strategy. - If granted, the task executes immediately.
Group rate limiting
Group rate limits share a limiter bucket across related tasks.
groupRateLimit: limiter policy for the shared group bucketgroupRateKey: optional static key (if omitted, Stem resolves from header)groupRateKeyHeader: header used whengroupRateKeyis not set (default:tenant)groupRateLimiterFailureMode(default:failOpen):failOpen: continue execution if limiter backend failsfailClosed: requeue/retry when limiter backend fails
class GroupRateLimitedTask extends TaskHandler<void> {
String get name => 'demo.groupRateLimited';
TaskOptions get options => const TaskOptions(
groupRateLimit: const RateLimit.perMinute(20),
groupRateKeyHeader: 'tenant',
groupRateLimiterFailureMode: RateLimiterFailureMode.failClosed,
maxRetries: 5,
);
Future<void> call(TaskContext context, Map<String, Object?> args) async {
final tenant = args['tenant'] as String? ?? 'global';
print('Handled group-rate-limited task for $tenant');
}
}
Redis-backed limiter example
The packages/stem/example/rate_limit_delay demo uses the shipped Redis
token-bucket limiter. It:
- shares tokens across multiple workers,
- uses Redis server time and an atomic Lua refill/acquire operation,
- reschedules denied tasks with retry metadata.
Observability
When a task is rate limited:
context.meta['rateLimited']is set on the retry attempt,taskRetrysignals include retry metadata,- worker logs show the limiter decision (if you log it).
Keying behavior
The worker uses a default rate-limit key of:
<taskName>:<tenant>
If no tenant header is set, it defaults to global. Add a tenant header when
enqueuing tasks to enforce per-tenant limits.
Redis limiter wiring
The rate_limit_delay example reads STEM_RATE_LIMIT_URL to point the limiter
at Redis. Use a dedicated Redis DB or key prefix to keep limiter state isolated
from your broker/result backend.
The Redis limiter is constructed with RedisRateLimiter.connect(...) from
stem_redis. stem_postgres provides the equivalent
PostgresRateLimiter.connect(...); it uses a server-clock token bucket with a
row lock inside one transaction.
Tips
- Use shared Redis for global limits across worker processes.
- Keep the rate limit key stable (by default it uses task name + tenant).
- Start with generous limits, then tighten after observing throughput.
Next steps
- See Tasks & Retries for other
TaskOptionsknobs. - Use Observability to instrument rate-limited flows.