In Celery, how do group, chunks() and starmap() differ when resizing 10,000 archived photos, in task messages sent and in parallelism?
answer
- one message per call, or one for all
- sequential inside one task
- starmap unpacks tuples
- a group of starmap batches
basics
~20 sIn Celery, a group sends one message per call, so 10,000 photos means 10,000 parallel tasks. starmap() sends one message and runs every call in sequence inside one task. chunks(it, n) sends a group of starmap tasks, each running n calls in sequence.
solid answer
~40 sA `group` of 10,000 `make_thumbnail` signatures publishes 10,000 messages; each call is its own task with its own id, state and result, spread across all workers. `make_thumbnail.map(ids)` and `make_thumbnail.starmap(pairs)` send **one** message to the built-in `celery.map` or `celery.starmap` task, which calls the function for every item in sequence inside a single worker: `map` passes each item as one argument, `starmap` unpacks each tuple into positional arguments. `make_thumbnail.chunks(pairs, 100)` sits in between: it splits the pairs into batches of 100 and sends a group of 100 `starmap` tasks, so batches run in parallel while calls within a batch run in sequence. Chunking cuts messaging overhead for many tiny calls; the cost is coarser failure, because one raising call fails its whole batch and the calls after it in that batch never run.
code
python · 20 linesfrom celery import Celery, group
app = Celery('photos', broker='redis://localhost:6379/0',
backend='redis://localhost:6379/1')
@app.task
def make_thumbnail(photo_id, size):
return f'thumbs/{photo_id}_{size}.jpg'
photo_ids = range(1, 10_001)
pairs = [(pid, 320) for pid in photo_ids]
# 10,000 messages, one task per photo
group(make_thumbnail.s(pid, 320) for pid in photo_ids).delay()
# 1 message; one worker runs all 10,000 calls in sequence
make_thumbnail.starmap(pairs).delay()
# 100 messages; 100 batches in parallel, 100 calls each
make_thumbnail.chunks(pairs, 100).apply_async()go deeper
Recall that a group sends one task per call, while map and starmap send one task that loops over every item.
Explain how chunks turns batches into a group of starmap tasks, how map and starmap pass arguments, and what the result looks like.
Show judgment on chunk size: throughput against per-item failure, redelivery cost and the loss of per-item retries, states and time limits.
Discuss when batching belongs in Celery's canvas versus inside a task or in a bulk-processing tool, given observability and failure-isolation needs.
## The problem: many small calls A photo-sharing app decides to generate a new 320 px thumbnail for 10,000 archived photos. Each resize takes a fraction of a second. The question is how to hand that work to Celery: as 10,000 separate tasks, as one task, or as something in between. Celery's canvas offers all three, and they differ in **messages sent**, **parallelism** and **failure granularity**. ## group: one message per call `group(make_thumbnail.s(pid, 320) for pid in photo_ids)` creates one signature per photo and sends them all. - **Messages:** 10,000 task messages, and 10,000 results if a result backend stores them. - **Parallelism:** maximal; any free worker slot takes the next call. - **Failure:** each call succeeds or fails on its own, with its own task id and state. The cost is overhead: publishing, routing, acknowledging and storing a result per call can outweigh a very short resize. ## map and starmap: one message for all `make_thumbnail.map(it)` and `make_thumbnail.starmap(it)` build signatures for Celery's built-in tasks **`celery.map`** and **`celery.starmap`**. Sending one publishes a **single message** carrying the whole list. The worker that receives it calls the task function for every item, one after another, in that one process: - `map` calls `make_thumbnail(item)` for each item: one argument per call. - `starmap` calls `make_thumbnail(*item)`: each item is a tuple unpacked into positional arguments, so `[(42, 320), (43, 320)]` becomes `make_thumbnail(42, 320)`, then `make_thumbnail(43, 320)`. The items are plain function calls inside one task, so there is **one task id, one state and one result**: a list of return values. There is **no parallelism**, and the first exception fails the whole task, leaving the remaining items unprocessed. ## chunks: batches in parallel `make_thumbnail.chunks(pairs, 100)` splits the iterable into batches of `n` and turns each batch into a `starmap` task. For 10,000 pairs and `n=100`: 1. The iterable is cut into 100 batches of 100 pairs. 2. Each batch becomes a `celery.starmap` signature. 3. The batches are sent as a **group**: 100 messages, run in parallel across workers. 4. Inside each batch, the 100 resizes run in sequence. Calling `.get()` on the result returns a list of lists, one inner list per batch. `chunks(...).group()` gives you the underlying group if you want to manipulate it before sending, for example to stagger start times with `skew()`. ## Side by side | Approach | Messages for 10,000 photos | Parallelism | A failing photo affects | |---|---|---|---| | `group` | 10,000 | Across all workers | Only its own task | | `starmap` / `map` | 1 | None; one worker, in sequence | The whole task; later items never run | | `chunks(it, 100)` | 100 | Across batches; sequential within one | Its batch; later items in that batch never run | ## Choosing - **Many tiny, uniform calls:** chunks. Celery's guide notes that chunking rarely hurts parallelism on a busy cluster and can raise throughput considerably by avoiding per-call messaging overhead. - **Few, slow or independent calls** that need their own retries, states or time limits: a group, so each photo is its own task. - **A short list** where one task is simpler to follow: starmap or map. - **Result storage follows the task count.** A group writes 10,000 small results to the result backend; each starmap task writes one list, so chunks with `n=100` writes 100 lists of 100 paths: fewer writes, larger values. - **Pick `n` for the failure blast radius as well as throughput.** With `n=1000`, one corrupt image stops up to 999 other resizes in that batch, and a redelivered batch redoes all of them. Resizes that can be safely repeated make that tolerable. - **Catch per-item errors inside the task** if one bad photo must not stop its batch: return a status for that photo instead of raising. ## Where the batches are sent from Calling a chunks signature, whether as `sig()` or `sig.apply_async()`, builds the group of `starmap` tasks in the **calling process** and publishes it from there. So a web request that chunks 10,000 photos pays for splitting the list and sending 100 messages before it returns; for very large inputs, send one small task that does the chunking on a worker instead.
- In a Celery chunks workflow, what happens when one photo in a batch raises?The batch is one `celery.starmap` task running a list comprehension over its items, so the exception fails that task at that item: the calls after it in the same batch never run, and the batch's state is `FAILURE`. Other batches are separate tasks and are unaffected. Catch per-photo errors inside `make_thumbnail` and return a status if one bad image must not stop its neighbours.
- How do you pick the chunk size n for a Celery chunks call?Balance messaging overhead against parallelism and blast radius. `n` should make each batch run long enough that publish, ack and result costs are small beside it, while leaving at least as many batches as worker slots so the cluster stays busy. Smaller `n` limits how much work one failure stops and how much a redelivered batch repeats.
saying these in an interview costs you the question
- map() and starmap() run their items in parallel across workers.
- map() unpacks each tuple into positional arguments.
- Chunking always makes the job slower because it removes parallelism.
- One failing photo in a chunk fails only that photo's call.
- chunks(it, n) splits the work into n batches rather than batches of n items.