Batch Enqueue
Insert many jobs in a single SQLite transaction with task.map() and enqueue_many(). Includes per-item result tracking for batched tasks.
Insert many jobs in a single SQLite transaction with task.map() and enqueue_many(). Includes per-item result tracking for batched tasks.
Insert many jobs in a single database transaction (SQLite or Postgres) for high throughput.
Batch enqueue (task.map() / enqueue_many()) writes many separate jobs
at once — one job per item. Task batching
(@queue.task(batch=…)) does the opposite: it collects many .delay() calls
into one job. Use batch enqueue for bulk-insert throughput; use task
batching to coalesce small units of work.
task.map()@queue.task()
def process(item_id):
return fetch_and_process(item_id)
# Enqueue 1000 jobs in one transaction
jobs = process.map([(i,) for i in range(1000)])queue.enqueue_many()# Basic batch — same options for all jobs
jobs = queue.enqueue_many(
task_name="myapp.process",
args_list=[(i,) for i in range(1000)],
priority=5,
queue="processing",
)
# Full parity with enqueue() — per-job overrides
jobs = queue.enqueue_many(
task_name="myapp.process",
args_list=[(i,) for i in range(100)],
delay=5.0, # uniform 5s delay for all
unique_keys=[f"item-{i}" for i in range(100)], # per-job dedup
metadata='{"source": "batch"}', # uniform metadata
expires=3600.0, # expire after 1 hour
result_ttl=600, # keep results for 10 minutes
)Per-job lists (delay_list, metadata_list, expires_list,
result_ttl_list) override uniform values when both are provided. See the
API reference for the full parameter list.
When using task-level batching (@queue.task(batch=...)), you can get an
individual result for each enqueued item rather than the whole batch list.
Enable it with per_item_results=True:
from flexiq import BatchItemResult
@queue.task(
batch={"max_size": 50, "max_wait_ms": 500, "per_item_results": True}
)
def process(items: list[int]) -> list[BatchItemResult]:
results = []
for i, item in enumerate(items):
try:
value = expensive_computation(item)
results.append(BatchItemResult.success(item_index=i, result=value))
except Exception as exc:
results.append(BatchItemResult.failure(item_index=i, error=str(exc)))
return resultsEach caller gets back a BatchedJobResult handle. Calling .result() blocks
until the batch flushes and returns this caller's per-item value:
h0 = process.delay(10)
h1 = process.delay(20)
h2 = process.delay(30)
# Each call returns its own value, not the whole list.
assert h0.result(timeout=15) == 100
assert h1.result(timeout=15) == 200
assert h2.result(timeout=15) == 300BatchItemResultfrom flexiq import BatchItemResult
# Convenience constructors:
item = BatchItemResult.success(item_index=2, result="done")
item = BatchItemResult.failure(item_index=0, error="network timeout")| Field | Type | Description |
|---|---|---|
status | "success" | "failure" | Outcome for this item |
result | Any | Return value (success only; None otherwise) |
error | str | None | Error message (failure only; None otherwise) |
item_index | int | Position in the original batch (≥ 0) |
When any item reports status="failure", the worker raises
BatchPartialFailureError and the whole batch retries (standard retry
policy). On retry, the task receives the same items again — write
the task function to be idempotent at the per-item level:
@queue.task(batch={"max_size": 50, "max_wait_ms": 500, "per_item_results": True})
def idempotent_process(items: list[int]) -> list[BatchItemResult]:
results = []
for i, item in enumerate(items):
# Check external state before acting — safe to replay.
if already_processed(item):
results.append(BatchItemResult.success(item_index=i, result=cached_result(item)))
else:
results.append(BatchItemResult.success(item_index=i, result=do_work(item)))
return resultspartial_failures()After a batch completes, call .partial_failures() to get the list of failed
items (empty list on full success):
handle = process.delay(42)
failures = handle.partial_failures(timeout=15)
# Returns [] on full success; list[BatchItemResult] otherwise.per_item_results=True and idempotent=True cannot be combined — the
decorator will raise ValueError at decoration time. Each item's identity
comes from item_index, not from a dedup key on the whole job.