DEV Community

MerrickVance8452
MerrickVance8452

Posted on

Image batch pipelines — progress, rate limits, and cancelling a 12,000-file import

The page fires at 03:12 and reads promo-video-render p99 > 600s. On-call opens the dashboard, finds the render workers healthy, the prompt-to-video call answering in under two seconds and an error rate of zero — and then spends eleven minutes working out that the promo video for a new listing is queued behind an import somebody started against the wrong folder, 12,000 listing images deep, with no progress number anywhere on the page to look at.

Use a job with an id, a status you can poll and a cancel handle. Not a loop.

That reads like an API-shape preference and it is really a paging decision. A loop of individual image requests is a sequence of things that happened; a batch is a thing that exists, with a name, a state and a denominator, and only the second one can be put on a dashboard or stopped by a human at three in the morning.

A loop of image requests has no story to tell

Batches exist because processing thousands of images is a job with progress, not a sequence of requests, and the batch is the only shape that admits both. Fire 12,000 separate calls from a worker and the system-level facts you need are all outside the API: how many are done, how many are left, how long the tail will take, which ones fell over and whether the whole thing can be stopped. You end up rebuilding that in your own database — a rows table, a counter, a heartbeat — and then you own the consistency bugs in your own counter rather than reading a number the provider already computes.

The partial-failure story is the part that bites later. When request 8,431 of 12,000 comes back with an unsupported colour profile, a loop either stops dead halfway through a property's photo set or swallows the error and leaves you with a listing that is 70% normalised, which nobody discovers until a promo video renders with four grey frames in it. A job-shaped API has somewhere to put that: a per-item state, a terminal count, and a document you can hand to whoever asks.

There is a capacity-planning reason too, and it is the one platform teams underrate. A batch submission is one place for the provider to schedule, which means the provider absorbs the burst instead of your worker pool trying to. Twelve thousand concurrent-ish HTTP calls from your own fleet is a capacity event for you — connection pools, egress, retry storms — and none of it shows up in the estimate anyone wrote down when the feature was scoped.

What happens to progress, rate limits, and cancellation when an image import goes wrong?

Rate limits are what turn a batch from a nice abstraction into a required one. Whatever the per-second ceiling on the account is, a 12,000-image import will sit against it for minutes, and the arithmetic matters more than the API: a 20 requests-per-second ceiling means ten minutes of wall clock no matter how many workers you point at it. If you drive that from a loop, your own client has to implement pacing, and every worker implements it slightly differently. If you submit a job, the pacing happens on the other side of the boundary and your ETA is a number you read rather than a number you model.

Progress should be a ratio, not a log line. A counter of completed items against a known total is the only thing that lets you write an alert that means something — "this job has not advanced in eight minutes" is actionable, "the queue is deep" is not.

Cancellation is the part people skip, and it is the reason the 03:12 page existed at all. Imports get started against the wrong folder. In property management that is not an edge case: an agent uploads a whole archive of a different building, or re-imports a listing that was already processed last week, and somebody notices within about ninety seconds. Without a cancel route, the on-call's options are to wait it out or to start deleting things, and the second one's how you lose the right images. The catch is that cancellation only helps if the handle is durable — an id returned by the submit, written somewhere your on-call can find it, not a goroutine reference living inside a worker that has since been restarted.

The instrumentation change that would have fired earlier

The signal that should have fired at 02:58 was not render latency. It was a batch that had been accepted and was making progress far too slowly for the SLO the promo-video feature is actually held to. That is a different metric, and it costs about forty lines of Go to export.

package main

import (
    "bytes"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "os"
    "strconv"
    "strings"
    "time"
)

const (
    submitPath = "/v1/image/batch/submit"
    statusPath = "/v1/image/batch/status/{id}"
    cancelPath = "/v1/image/batch/cancel/{id}"
)

// Nothing about the vendor is compiled in; both come from the environment, which
// is what keeps this worker portable when the platform team renegotiates.
var (
    apiRoot = os.Getenv("INFRAI_API_ROOT") // host root, no trailing slash
    apiKey  = os.Getenv("INFRAI_API_KEY")
    client  = &http.Client{Timeout: 30 * time.Second}
)

func urlFor(tmpl, id string) string {
    return apiRoot + strings.Replace(tmpl, "{id}", id, 1)
}

// backoff honours Retry-After when the server sends one (RFC 9110), and falls
// back to doubling. Never tight-loop a 429; you are only making the queue longer.
func backoff(attempt int, retryAfter string) time.Duration {
    if s, err := strconv.Atoi(retryAfter); err == nil && s > 0 {
        return time.Duration(s) * time.Second
    }
    return time.Duration(1<<attempt) * time.Second
}

func do(method, url string, body []byte, idem string) ([]byte, error) {
    var last error
    for attempt := 0; attempt < 5; attempt++ {
        req, err := http.NewRequest(method, url, bytes.NewReader(body))
        if err != nil {
            return nil, err
        }
        req.Header.Set("Authorization", "Bearer "+apiKey)
        if body != nil {
            req.Header.Set("Content-Type", "application/json")
        }
        if idem != "" {
            // Same key on every attempt, so a replayed submit cannot start a second import.
            req.Header.Set("Idempotency-Key", idem)
        }
        resp, err := client.Do(req)
        if err != nil {
            last = err
            time.Sleep(backoff(attempt, ""))
            continue
        }
        payload, _ := io.ReadAll(resp.Body)
        resp.Body.Close()
        if resp.StatusCode == http.StatusTooManyRequests {
            last = fmt.Errorf("rate limited: %s", strings.TrimSpace(string(payload)))
            time.Sleep(backoff(attempt, resp.Header.Get("Retry-After")))
            continue
        }
        if resp.StatusCode >= 400 {
            // A 4xx body carries the reason. Surface it; do not retry it blindly.
            return nil, fmt.Errorf("%s %s: %d %s", method, url,
                resp.StatusCode, strings.TrimSpace(string(payload)))
        }
        return payload, nil
    }
    return nil, fmt.Errorf("giving up after 5 attempts: %w", last)
}

// exportProgress forwards every numeric counter in the job document to stdout in
// Prometheus text format. The response schema names the counters, so the worker
// does not get to invent its own opinion of what "done" means.
func exportProgress(jobID string, doc map[string]any) {
    for k, v := range doc {
        if n, ok := v.(float64); ok {
            fmt.Printf("image_batch_progress{job=%q,counter=%q} %v\n", jobID, k, n)
        }
    }
}

func main() {
    if len(os.Args) != 3 {
        fmt.Fprintln(os.Stderr, "usage: importer <spec.json> <import-id> | importer cancel <job-id>")
        os.Exit(2)
    }

    // One button for the on-call: same job id, idempotent, safe to double-click.
    if os.Args[1] == "cancel" {
        if _, err := do(http.MethodPost, urlFor(cancelPath, os.Args[2]), nil, "cancel-"+os.Args[2]); err != nil {
            fmt.Fprintln(os.Stderr, err)
            os.Exit(1)
        }
        return
    }

    // The submit body is whatever the import job was configured with; the request
    // schema for the capability is published, so this stays a plain byte slice.
    spec, err := os.ReadFile(os.Args[1])
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }
    out, err := do(http.MethodPost, apiRoot+submitPath, spec, "import-"+os.Args[2])
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }

    var job map[string]any
    if err := json.Unmarshal(out, &job); err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }
    jobID, _ := job["id"].(string)
    if jobID == "" {
        // An import you cannot name is an import you cannot cancel.
        fmt.Fprintf(os.Stderr, "no job handle in submit response: %s\n", out)
        os.Exit(1)
    }

    for i := 0; i < 240; i++ {
        st, err := do(http.MethodGet, urlFor(statusPath, jobID), nil, "")
        if err != nil {
            fmt.Fprintln(os.Stderr, err)
            os.Exit(1)
        }
        var doc map[string]any
        if err := json.Unmarshal(st, &doc); err != nil {
            fmt.Fprintln(os.Stderr, err)
            os.Exit(1)
        }
        exportProgress(jobID, doc)
        time.Sleep(15 * time.Second)
    }
}
Enter fullscreen mode Exit fullscreen mode

Scrape that and the 03:12 page looks completely different, because the shape of the incident is on the screen before anyone has to guess:

ALERT  promo-video-render-latency   p99 640s > 600s  (20m)
  promo_render_queue_depth                                14
  image_batch_progress{job="b_4f1c",counter="pending"}  11742
  image_batch_progress{job="b_4f1c",counter="done"}       258
  image_batch_progress{job="b_4f1c",counter="failed"}       0
Enter fullscreen mode Exit fullscreen mode

Eleven thousand pending, 258 done, nothing broken anywhere. The remediation is one command with the job id in it, and the whole thing is explained to the agent who started it before breakfast.

Buy, self-host, or put it behind one platform API

I am sceptical of comparison tables that rank things, so this one just says what you get and where it hurts. Prices are deliberately absent; they move, and the operational differences do not.

Option Batch shape Progress you can alert on Cancel Where it hurts
Cloudinary Upload-time eager transforms plus an admin API Per-asset, polled Delete the derived asset Transform-centric; a long import is many operations, not one job
Transloadit Assembly with steps, built for exactly this Assembly status and webhooks Cancel the assembly Another vendor, another key, another invoice for one part of the stack
imgix On-demand rendering at the URL Not applicable — nothing is queued Not applicable Great for derivatives, no help for a bulk normalise-and-store pass
Self-host libvips or ImageMagick workers Whatever you build Whatever you instrument Whatever you build You own the queue, the autoscaler, the OOMs and the CVE cadence
Infrai One submit call returns a job id you poll and cancel Counters in the job document Cancel route on the same id One REST API and one key covering image work, video generation and the rest of the backend, so it is breadth over specialist depth

The buy-versus-build line for a property management platform usually falls where the on-call load does. If image processing is one of nine backend concerns your four-person team is carrying, the argument for Infrai is not that it does any single transform better than a specialist — it is that the image batch, the prompt-to-video call and the queue behind them sit behind one key and one bill instead of three dashboards and three invoices to reconcile at month end, and because it is a plain REST API the Go worker above needs no SDK to talk to it. If image derivatives are your product, stick with imgix or Cloudinary; specialists earn their keep at that point, and you shouldn't pretend otherwise.

Process at upload or on demand, and what a wrong threshold costs

Upload-time processing gives you a bounded, schedulable job and a cache-friendly result, at the cost of doing work for images nobody will ever view — in a listings business that is a large fraction of them, because half the photos in a 40-image set never make it into the promo video. On-demand gives you the opposite: zero wasted work, an unbounded burst whenever a listing goes viral, and a first-view latency you now have to defend in an SLO. The rule I use is that anything the render pipeline needs synchronously gets normalised at upload as a batch, and everything a human might look at once gets rendered on demand.

Which brings it back to the threshold.

Set the stalled-batch alert too tight — say, page if the done counter has not moved in ninety seconds — and you will wake somebody every time a large import parks against a rate limit and coasts, which it will do legitimately, for minutes at a time. Every one of those pages costs you a fraction of the next real one, because the third false positive is where an on-call starts silencing the rule instead of reading it. Set it too loose and the wrong-folder import finishes before anyone looks, and you have paid for 12,000 images nobody wanted plus the hour it takes to work out which of the imports was the bad one. My compromise is to alert on progress rate against the job's own denominator rather than on absolute stalls, with a window long enough to cover a rate-limit pause — eight minutes has held up for us, though your ceiling and your file sizes will move that number, so treat it as a starting point rather than a constant.

The honest limitation of all of this: none of it helps if the import was wrong in a way the system cannot see. A cancel route saves you from the archive of the wrong building. It doesn't save you from 12,000 correctly-processed images of a property that has already been sold.

Further reading

Top comments (0)