Worker Fleet and Scaling

Worker Fleet and Scaling An architecture diagram generated by Archify. apps/mail · adds jobs · Architecture component apps/mail adds jobs Redis queue · BullMQ · Architecture component Redis queue BullMQ Listener · task 1 · ECS service · apps/worker › task 1 Listener task 1 Child pool · C processes · ECS service · apps/worker › task 1 Child pool C processes Listener · task 2 · ECS service · apps/worker › task 2 Listener task 2 Child pool · C processes · ECS service · apps/worker › task 2 Child pool C processes Listener · task N, starting · ECS service · apps/worker › task N (new) Listener task N, starting Sweeper · repeatable job · Architecture component Sweeper repeatable job CloudWatch · backlog per task · Architecture component CloudWatch backlog per task Auto Scaling · target tracking · Architecture component Auto Scaling target tracking Postgres · rows · pipeline tables · Architecture component Postgres rows · pipeline tables Blob storage · raw/<sha>.eml · Architecture component Blob storage raw/<sha>.eml add job job 1 job 2 run run re-add missed backlog over target start task commit rows write MIME row + artifacts raw MIME ECS service · apps/worker task 1 task 2 task N (new)

Pulling work

  • • Every task's listener pulls from one Redis queue
  • • Each claim is atomic and locked: one listener per job
  • • A listener runs up to C jobs in pooled child processes
  • • Every pool reads rows + MIME and saves artifacts (drawn for task 1)

Scaling out

  • • The sweeper publishes backlog per task every minute
  • • Over target, Auto Scaling raises the ECS desired count
  • • ECS starts task N; its listener starts pulling at once

Scaling in

  • • Auto Scaling lowers the count once the queue drains
  • • ECS sends SIGTERM; worker.close() waits for running jobs
  • • A job cut off after 120 s is retried on another task