Machine Coding Problem

Thumbnail Service

macoAllmediaasync-job-processing
Commonly Asked By:PinterestNetflixDropbox

Requirements & System Scope

Functional Scope (In-Scope)

  • Asynchronous Job Queueing: Offloads image transcoding from the API gateway using an in-memory Blocking Queue and worker pool.
  • Idempotent Operations: Uses hashes computed from original Image IDs and Target Resolutions to reuse existing jobs.
  • Robust Exponential Back-off Retries: Automatically re-schedules failed sub-jobs with adaptive back-off.
  • Atomic State Updates & Event callbacks: Maintains strict thread safety and pushes callbacks to webhook endpoints on completion.

Explicit Boundaries (Out-of-Scope)

  • Physical Image Resizing & Compression: Mocks pixel rendering pipelines and image decoders (like libjpeg).
  • Webhooks / Dynamic Load Balancing: Simple console webhook logging replaces fully qualified web hook endpoints.

Class Diagram & Entity Relationships

Worker-pool thread allocations, non-blocking queues, and job states:

Loading...
  • ThumbnailJob Entity: Tracks execution variables, back-off counts, resolution sets, and destination paths.
  • BlockingQueue Pool: Decouples high-volume API requests from active thumbnail generation threads.

Design Patterns & SOLID Principles

  • Strategy Pattern (Image Processing Strategy): Decouples the concrete image resizing/transcoding algorithm (e.g. HighQualityImageProcessor with Lanczos interpolation vs FastImageProcessor for high-throughput scaling) from the worker execution pool, facilitating runtime swappability and adhering to OCP.
  • Producer-Consumer Pattern: Drives task processing using concurrent threads consuming jobs from a shared blocking queue.
  • Idempotent Consumer Pattern: Prevents duplicating resource-intensive image compression tasks for the same requested dimensions.
  • Single Responsibility Principle (SRP): Decouples file uploads, job queueing, image resizing, and completion webhooks into independent classes.

Core Execution Workflows

Stateless Retry Delays & Queue Workers

  1. Idempotency Verification & Registration:
    1. Receive image uploading events containing required target output dimensions.
    2. Compute a SHA-256 hash of the image metadata. If the hash exists and matches an active job, link the callback to the existing job and return.
    3. Otherwise, create a new pending job, insert it into the tracking index, and push the ID into the processing queue.
  2. Exponential Back-off Rescheduling:
    1. If image buffer resizing fails due to transient faults, increment the retry counter.
    2. Verify if the retry count is under the maximum limit. If yes, calculate a delayed schedule: delay = 2 ^ retries * baseline.
    3. Re-queue the job after the calculated delay. If the retry limit is exceeded, update the job status to FAILED and invoke error callbacks.

Concurrency & Thread Safety Strategy

Ensuring reliable processing under heavy load:

  • Concurrent Collections: Stores active job metadata and callbacks in concurrent hash maps to prevent race conditions during updates.
  • Safe Queue Synchronization & Locks: Utilizes thread-safe blocking queues and Locks in Python to balance load across resizer workers without requiring manual synchronization.

Complete Clean Code Blueprint

Production reference implementations demonstrating thread pools, blocking queues, back-off retry logic, and webhook callbacks in Java and Python:

// โ”€โ”€โ”€ JAVA BLUEPRINT โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

enum JobStatus {
    PENDING, PROCESSING, DONE, FAILED
}

class Resolution {
    private final int width;
    private final int height;

    public Resolution(int width, int height) {
        this.width = width;
        this.height = height;
    }

    public int getWidth() { return width; }
    public int getHeight() { return height; }

    @Override
    public String toString() {
        return width + "x" + height;
    }
}

class ThumbnailJob {
    private final String id;
    private final String originalImageId;
    private final List<Resolution> resolutions;
    private JobStatus status;
    private int retries;
    private final int maxRetries;
    private final Map<String, String> generatedPaths; // resolution -> storagePath
    private String errorMessage;
    private final long createdAtMs;

    public ThumbnailJob(String id, String originalImageId, List<Resolution> resolutions, int maxRetries) {
        this.id = id;
        this.originalImageId = originalImageId;
        this.resolutions = new ArrayList<>(resolutions);
        this.status = JobStatus.PENDING;
        this.retries = 0;
        this.maxRetries = maxRetries;
        this.generatedPaths = new ConcurrentHashMap<>();
        this.createdAtMs = System.currentTimeMillis();
    }

    public String getId() { return id; }
    public String getOriginalImageId() { return originalImageId; }
    public List<Resolution> getResolutions() { return resolutions; }
    public synchronized JobStatus getStatus() { return status; }
    public synchronized void setStatus(JobStatus status) { this.status = status; }
    public synchronized int getRetries() { return retries; }
    public synchronized void incrementRetries() { this.retries++; }
    public int getMaxRetries() { return maxRetries; }
    public Map<String, String> getGeneratedPaths() { return generatedPaths; }
    public synchronized String getErrorMessage() { return errorMessage; }
    public synchronized void setErrorMessage(String msg) { this.errorMessage = msg; }
}

interface ImageProcessor {
    String process(String imageId, Resolution resolution) throws Exception;
}

class HighQualityImageProcessor implements ImageProcessor {
    @Override
    public String process(String imageId, Resolution res) throws Exception {
        // Simulating unstable imaging network
        if (Math.random() < 0.15) {
            throw new RuntimeException("Transcoder failed to load image buffers.");
        }
        return "/storage/thumbnails/hq_" + imageId + "_" + res.getWidth() + "x" + res.getHeight() + ".jpg";
    }
}

class FastImageProcessor implements ImageProcessor {
    @Override
    public String process(String imageId, Resolution res) throws Exception {
        return "/storage/thumbnails/fast_" + imageId + "_" + res.getWidth() + "x" + res.getHeight() + ".jpg";
    }
}

class ThumbnailService {
    private final ConcurrentHashMap<String, ThumbnailJob> jobs = new ConcurrentHashMap<>();
    private final BlockingQueue<String> jobQueue = new LinkedBlockingQueue<>();
    private final ConcurrentHashMap<String, List<String>> callbacks = new ConcurrentHashMap<>();
    private final ExecutorService workerPool;
    private final ScheduledExecutorService retryScheduler = Executors.newScheduledThreadPool(1);
    private final int defaultMaxRetries = 3;
    private final ImageProcessor imageProcessor;

    public ThumbnailService(int threadCount, ImageProcessor processor) {
        this.imageProcessor = Objects.requireNonNull(processor);
        this.workerPool = Executors.newFixedThreadPool(threadCount);
        for (int i = 0; i < threadCount; i++) {
            workerPool.submit(this::processQueue);
        }
    }

    // Submit Job with Idempotency Validation
    public String submitJob(String imageId, List<Resolution> resolutions, String callbackUrl) {
        String jobId = calculateIdempotencyKey(imageId, resolutions);

        // Register callback url if provided
        if (callbackUrl != null) {
            callbacks.computeIfAbsent(jobId, k -> new CopyOnWriteArrayList<>()).add(callbackUrl);
        }

        // Idempotency Check
        ThumbnailJob existing = jobs.get(jobId);
        if (existing != null) {
            if (existing.getStatus() == JobStatus.FAILED) {
                // Retry/reset failed job
                existing.setStatus(JobStatus.PENDING);
                existing.setErrorMessage(null);
                jobQueue.offer(jobId);
            }
            return jobId;
        }

        ThumbnailJob newJob = new ThumbnailJob(jobId, imageId, resolutions, defaultMaxRetries);
        jobs.put(jobId, newJob);
        jobQueue.offer(jobId);

        return jobId;
    }

    public ThumbnailJob getJobStatus(String jobId) {
        return jobs.get(jobId);
    }

    private void processQueue() {
        try {
            while (!Thread.currentThread().isInterrupted()) {
                String jobId = jobQueue.take(); // Blocks until job is ready
                ThumbnailJob job = jobs.get(jobId);
                if (job == null) continue;

                processJob(job);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    private void processJob(ThumbnailJob job) {
        job.setStatus(JobStatus.PROCESSING);

        try {
            for (Resolution res : job.getResolutions()) {
                String resultPath = imageProcessor.process(job.getOriginalImageId(), res);
                job.getGeneratedPaths().put(res.toString(), resultPath);
            }
            job.setStatus(JobStatus.DONE);
            triggerCallbacks(job);
        } catch (Exception e) {
            job.incrementRetries();
            if (job.getRetries() <= job.getMaxRetries()) {
                long backoffDelayMs = (long) Math.pow(2, job.getRetries()) * 100L;
                job.setStatus(JobStatus.PENDING);
                job.setErrorMessage("Retrying. Error: " + e.getMessage());
                
                retryScheduler.schedule(() -> jobQueue.offer(job.getId()), backoffDelayMs, TimeUnit.MILLISECONDS);
            } else {
                job.setStatus(JobStatus.FAILED);
                job.setErrorMessage("Max retries reached. Root cause: " + e.getMessage());
                triggerCallbacks(job);
            }
        }
    }

    private void triggerCallbacks(ThumbnailJob job) {
        List<String> urls = callbacks.get(job.getId());
        if (urls == null) return;
        
        for (String url : urls) {
            System.out.println("NOTIFY CALLBACK URL (" + url + ") -> Job: " + job.getId() + " Status: " + job.getStatus());
        }
    }

    private String calculateIdempotencyKey(String imageId, List<Resolution> resolutions) {
        List<Resolution> sorted = new ArrayList<>(resolutions);
        sorted.sort(Comparator.comparingInt(Resolution::getWidth).thenComparingInt(Resolution::getHeight));
        
        StringBuilder sb = new StringBuilder(imageId);
        for (Resolution res : sorted) {
            sb.append("|").append(res.toString());
        }

        try {
            MessageDigest digest = MessageDigest.getInstance("SHA-256");
            byte[] hash = digest.digest(sb.toString().getBytes(StandardCharsets.UTF_8));
            StringBuilder hexString = new StringBuilder();
            for (byte b : hash) {
                String hex = Integer.toHexString(0xff & b);
                if (hex.length() == 1) hexString.append('0');
                hexString.append(hex);
            }
            return hexString.toString().substring(0, 16);
        } catch (Exception e) {
            return String.valueOf(sb.toString().hashCode());
        }
    }

    public void shutdown() {
        workerPool.shutdownNow();
        retryScheduler.shutdownNow();
    }
}

public class Main {
    public static void main(String[] args) throws Exception {
        System.out.println("=== JAVA THUMBNAIL SERVICE DEMO ===");
        ThumbnailService service = new ThumbnailService(2, new HighQualityImageProcessor());

        List<Resolution> resolutions = Arrays.asList(
            new Resolution(150, 150),
            new Resolution(300, 300),
            new Resolution(600, 600)
        );

        String imageId = "profile_pic_2026";
        String callbackUrl = "https://myapi.com/webhooks/thumbnails";

        System.out.println("Submitting image processing job...");
        String jobId = service.submitJob(imageId, resolutions, callbackUrl);
        System.out.println("Job submitted. Generated Job ID: " + jobId);

        // Retrieve job status
        ThumbnailJob job = service.getJobStatus(jobId);
        System.out.println("Initial Job Status: " + job.getStatus());

        // Wait a bit for processing to progress (mocking asynchronous delay)
        System.out.println("Waiting for asynchronous workers to process...");
        Thread.sleep(1500);

        ThumbnailJob finalJobState = service.getJobStatus(jobId);
        System.out.println("Final Job Status: " + finalJobState.getStatus());
        System.out.println("Generated paths: " + finalJobState.getGeneratedPaths());
        if (finalJobState.getErrorMessage() != null) {
            System.out.println("Errors encountered: " + finalJobState.getErrorMessage());
        }

        // Test idempotency
        System.out.println("Submitting identical job to test idempotency...");
        String duplicateJobId = service.submitJob(imageId, resolutions, callbackUrl);
        System.out.println("Duplicate request Job ID: " + duplicateJobId + " (Matches original: " + duplicateJobId.equals(jobId) + ")");

        service.shutdown();
        System.out.println("=== END OF DEMO ===");
    }
}

๐Ÿ’ฌReview

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...