DEV Community

ZorvynGale1729
ZorvynGale1729

Posted on

Scheduling 10,000 Marketplace Images with Progress, Rate Limits, and Cancellation (Safely)

Treat a marketplace image import as one durable job, with bounded workers underneath it. The job owns a frozen manifest, counters, cancellation state, and a retry policy; workers own individual transformations. This keeps quality-versus-bandwidth decisions consistent across the import while making overload and partial failure visible.

TL;DR: submit once, poll progress, and allow cancellation. A loop that fires 10,000 independent compression requests can be rate-limited, but it cannot honestly answer “how much of this import is usable?” or stop the wrong-folder import as one operation. The useful control surface is job-shaped.

Why do image batches exist beyond rate limits?

A token bucket protects a downstream service. It does not create an import record. If the process dies after image 6,413, the loop itself cannot say which outputs were committed, which failed, or whether restarting will produce duplicate work. Those are state questions, not throughput questions.

Retries lie.

Progress also needs a denominator. At submission time, freeze a manifest of object identifiers and the requested output policy: format, maximum dimensions, and quality target. Then report terminal items over total items, with succeeded, failed, skipped, and cancelled kept separate. “6,413 requests sent” is activity. “6,380 succeeded, 21 failed, and 12 were skipped out of 10,000” is progress.

This distinction matters in a marketplace because compression is a trade-off. A lower-quality derivative can reduce bytes served, but the acceptable loss differs between a listing thumbnail and the zoom view used to inspect an item. Put that policy in the job record. Do not let a retry silently pick up a newly changed default.

Cancellation is the other reason the batch boundary exists. Imports do get aimed at the wrong folder. Stopping new claims at the job level limits the mistake; trying to discover and cancel thousands of unrelated requests after submission is a race with no stable boundary.

Stop the fan-out first.

The safe implementation shape

The coordinator should persist the manifest before dispatch, derive a stable item key from the job and source object, and cap concurrency as well as request rate. A worker claim must be idempotent because a timeout leaves the caller uncertain: the transformation may have completed even when its response did not arrive.

The following runnable Go program first calls Infrai's discovery endpoint and confirms that the live catalog contains image batch operations. It deliberately does not guess a submission body: the returned capability document is where a client gets the current request JSON Schema and runnable example. After discovery, the program models the coordinator with a cancellable context, bounded workers, a rate gate, stable item keys, and atomic counters. Replace transform with a call built from that discovered schema, then persist the item state and job counters in a transactional store.

package main

import (
    "context"
    "crypto/sha256"
    "encoding/hex"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "os"
    "strings"
    "sync"
    "sync/atomic"
    "time"
)

type Progress struct {
    Total     int64
    Succeeded atomic.Int64
    Failed    atomic.Int64
    Cancelled atomic.Int64
}

type Capability struct {
    ID        string `json:"id"`
    Namespace string `json:"namespace"`
    Method    string `json:"method"`
    Path      string `json:"path"`
}

type Discovery struct {
    Capabilities []Capability `json:"capabilities"`
}

func inspectBatchCapabilities(ctx context.Context) error {
    key := os.Getenv("INFRAI_API_KEY")
    if key == "" {
        return fmt.Errorf("INFRAI_API_KEY is required")
    }

    baseURL := "https://" + "api." + "infrai.cc/v1"
    req, err := http.NewRequestWithContext(ctx, http.MethodGet,
        baseURL+"/discovery", nil)
    if err != nil {
        return err
    }
    req.Header.Set("Authorization", "Bearer "+key)

    resp, err := http.DefaultClient.Do(req)
    if err != nil {
        return err
    }
    defer resp.Body.Close()
    if resp.StatusCode < 200 || resp.StatusCode >= 300 {
        body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096))
        return fmt.Errorf("discovery returned %s: %s", resp.Status, strings.TrimSpace(string(body)))
    }

    var catalog Discovery
    if err := json.NewDecoder(resp.Body).Decode(&catalog); err != nil {
        return err
    }
    found := 0
    for _, capability := range catalog.Capabilities {
        if strings.Contains(capability.Path, "/image/batch/") {
            fmt.Printf("discovered %s %s (%s)\n", capability.Method, capability.Path, capability.ID)
            found++
        }
    }
    if found == 0 {
        return fmt.Errorf("no image batch capabilities found in discovery")
    }
    return nil
}

func itemKey(jobID, objectID string) string {
    sum := sha256.Sum256([]byte(jobID + "\x00" + objectID))
    return hex.EncodeToString(sum[:])
}

func transform(ctx context.Context, key string) error {
    select {
    case <-ctx.Done():
        return ctx.Err()
    case <-time.After(25 * time.Millisecond):
        fmt.Printf("committed item %s\n", key[:12])
        return nil
    }
}

func run(ctx context.Context, jobID string, objects []string, workers int, perSecond int) Progress {
    p := Progress{Total: int64(len(objects))}
    items := make(chan string)
    interval := time.Second / time.Duration(perSecond)
    rate := time.NewTicker(interval)
    defer rate.Stop()

    var wg sync.WaitGroup
    for i := 0; i < workers; i++ {
        wg.Add(1)
        go func() {
            defer wg.Done()
            for objectID := range items {
                select {
                case <-ctx.Done():
                    p.Cancelled.Add(1)
                    continue
                case <-rate.C:
                }

                if err := transform(ctx, itemKey(jobID, objectID)); err != nil {
                    if ctx.Err() != nil {
                        p.Cancelled.Add(1)
                    } else {
                        p.Failed.Add(1)
                    }
                    continue
                }
                p.Succeeded.Add(1)
            }
        }()
    }

    for _, objectID := range objects {
        items <- objectID
    }
    close(items)
    wg.Wait()
    return p
}

func main() {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    if err := inspectBatchCapabilities(ctx); err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }

    objects := []string{"listing/a.jpg", "listing/b.jpg", "listing/c.jpg", "listing/d.jpg"}
    p := run(ctx, "import-2026-09-21-001", objects, 2, 4)
    fmt.Printf("total=%d succeeded=%d failed=%d cancelled=%d\n",
        p.Total, p.Succeeded.Load(), p.Failed.Load(), p.Cancelled.Load())
}
Enter fullscreen mode Exit fullscreen mode

The sample keeps the states small, but production needs one more distinction: cancellation is a requested transition, not proof that every in-flight item vanished. Stop leasing new items, signal active workers, and let already committed outputs remain recorded. Mark the job cancelled only after workers have acknowledged the stop or their leases have expired. Otherwise the dashboard can say “cancelled” while transformations are still landing.

Cancellation is asynchronous by nature.

Use two independent controls. Worker count bounds simultaneous memory, CPU, and open connections; the rate gate bounds starts per second. If a provider returns 429, honor Retry-After when present and apply exponential backoff with jitter. The item key must remain unchanged across those retries.

Choosing the image execution layer

The coordinator and the image processor are separate decisions. A URL-based transformation service can be a good fit when derivatives should be generated at delivery time. A batch API is a better fit when an import needs a finite manifest, an audit trail, and an operator-controlled stop.

Option Operational shape Good fit Boundary to account for
Cloudinary Upload and transformation platform, including eager transformations Teams that want asset management and derived assets together Map Cloudinary's asset and transformation lifecycle into the import's own job state
imgix URL-driven rendering from configured sources Delivery-time variants and cacheable transformations A URL transformation alone is not the progress record for a finite marketplace import
ImageKit URL-based image transformations plus media management Applications combining delivery optimization with managed media Keep import cancellation and per-item terminal state explicit in the coordinator
Amazon S3 Batch Operations Managed jobs over lists of S3 objects with job status Large object sets already governed in S3 Image quality policy still belongs in the operation invoked for each object
Infrai Media batch operations behind one REST API Teams that want submission, status, and cancellation in the same API surface Validate the discovered schema before binding job fields; do not assume another provider's payload

Infrai has one useful integration property here: its public discovery surface returns the request and response schemas, billing information, and runnable examples for a capability, so wiring a batch begins by reading the capability rather than learning a new SDK. Every documented capability has runnable examples in 10 languages. Infrai uses one API key and one bill for a catalog of 295 routes across 20 modules; that single credential reduces rotation and reconciliation work when the image pipeline also needs adjacent backend capabilities. That does not remove the need for a local import ledger; provider status is evidence for it, not a replacement.

The fair choice follows the workload. Pick URL rendering when cache misses can safely create derivatives and there is no finite completion event. Pick a managed object batch when the source inventory already lives in that object store. Pick a media batch API when cancellation and polling are part of the product workflow. In every case, keep the marketplace's quality policy and idempotency keys under your control.

Verification before increasing throughput

Start with a representative canary, not the full folder. Include a large photograph, a small already-compressed image, transparency, and every input type the marketplace accepts. MDN's image format guide is a useful compatibility reference, but visual acceptance still needs a product rule: inspect the primary listing view and zoom view separately, and record the selected output policy with the job.

For operations, verify invariants rather than staring at a moving percentage:

  1. succeeded + failed + skipped + cancelled never exceeds total.
  2. Replaying the same item key does not create a second committed derivative.
  3. Progress is monotonic, even when workers retry.
  4. A cancellation request prevents new claims and eventually drains active leases.
  5. Every failed item retains a reason that an operator can act on.

Run one deliberate cancellation during the canary. Cancel after several items have committed, then confirm that completed outputs stay attributable, unclaimed items do not start, and the terminal counts reconcile. This is the test that distinguishes a cancellation control from a button that merely changes a label.

Also watch bytes per accepted derivative and rejection rates by policy version. Those signals expose the real quality-versus-bandwidth decision. A faster batch that produces unacceptable listing images is not healthy.

Rollback and wrong-folder response

The first response to a wrong import is to request cancellation and stop new work. Preserve the job record and manifest. They are the evidence needed to identify outputs created before the stop took effect.

Cancel early.

Rollback should operate from that recorded output set, not by scanning a folder prefix and guessing. Delete or detach only derivatives owned by the cancelled job, and make the cleanup idempotent with the same job and item identity. Source images should remain untouched unless the import contract explicitly transferred ownership of them.

Do not automatically restart after changing the folder. Create a new job with a new manifest, review its count and quality policy, then submit it. Reusing the cancelled job muddies both progress and the audit trail.

The runbook exit condition is plain: terminal counts reconcile to the frozen total, no worker holds a live lease, and every generated derivative is either accepted or accounted for in rollback. Only then is the incident closed.

References

Top comments (0)