Guide

Batches

Enqueue a group of jobs and run a callback when the whole batch completes.

Create and enqueue a batch#

A batch tracks a group of related jobs. Enqueue the jobs inside batch.enqueue and each is tagged with the batch id:

batch = Pgbus::Batch.new(
  on_finish: BatchFinishedJob,
  on_success: BatchSucceededJob,
  on_failure: BatchFailedJob,
  description: "Import users",
  properties: { initiated_by: current_user.id }
)

batch.enqueue do
  users.each { |user| ImportUserJob.perform_later(user.id) }
end

Open batches#

A batch stays open until it finishes. Call enqueue again to add another stage — total_jobs grows and the callbacks wait for the new jobs too:

batch = Pgbus::Batch.new(on_finish: BatchFinishedJob)
batch.enqueue { ExtractJob.perform_later }
batch.enqueue { TransformJob.perform_later }  # same batch, total_jobs == 2

A job running inside a batch reaches its own batch through batch and can add siblings the same way. Membership stays explicit: only jobs enqueued inside an enqueue block join the batch, so a fan-out from a batched job does not silently extend it.

app/jobs/extract_job.rb
class ExtractJob < ApplicationJob
  def perform
    rows = extract
    batch.enqueue do
      rows.each { |row| TransformJob.perform_later(row.id) }
    end
  end
end

Pgbus::Batch.find(batch_id) returns the same handle from anywhere — with description, properties, status, total_jobs, completed_jobs, failed_jobs, pending_jobs, progress_percentage and finished?. Adding to a batch that has already finished raises Pgbus::Batch::AlreadyFinished — at perform_later, before the job is sent, even if the handle you hold is stale.

Breaking change (pre-1.0): Pgbus::Batch.find used to return the raw attributes Hash. Read the values off the handle, or query Pgbus::BatchEntry directly for a row.

Callbacks#

CallbackFired when
on_finishThe batch finished (no outstanding execution rows remain), including after a dispatcher sweep repair.
on_successThe batch finished with zero failed jobs.
on_failureThe batch finished with at least one dead-lettered job. (`on_discard:` is a deprecated alias until 1.0.)

A callback job receives the batch properties hash as its argument:

app/jobs/batch_finished_job.rb
class BatchFinishedJob < ApplicationJob
  def perform(properties)
    user = User.find(properties["initiated_by"])
    ImportMailer.complete(user).deliver_later
  end
end

Configured callback jobs#

A callback can be a configured ActiveJob instance instead of a bare class, so it runs on the queue — and with the delay — you choose:

Pgbus::Batch.new(
  on_finish: BatchFinishedJob.new.set(queue: :critical, wait: 5.minutes)
)

.set options resolve when the batch is created, and the serialized job is stored on the batch row. At fire time it is enqueued on its configured queue with callback_batch_id pointing at the finished batch, so the callback reads the batch through batch rather than through a properties argument:

app/jobs/batch_finished_job.rb
class BatchFinishedJob < ApplicationJob
  def perform(*)
    Rails.logger.info "#{batch.completed_jobs}/#{batch.total_jobs} done"
    User.find(batch.properties["initiated_by"])
  end
end

A callback is never a member of the batch it reports on — its own batch_id is nil, so enqueueing it can never keep the batch open.

Existing installs need rails generate pgbus:add_batch_callback_jobs (or pgbus:update) for the jsonb callback columns. Bare callback classes keep the perform_later(properties) signature, deprecated at 1.0.

How batches work#

  1. Batch.new(...) creates a row in pgbus_batches with status: "pending".
  2. batch.enqueue { ... } tags each enqueued job with the batch id and, in one transaction before the message is sent, increments total_jobs and inserts a pgbus_batch_executions row (identity is the ActiveJob job_id). The increment is guarded on an unfinished batch — that guard is what raises AlreadyFinished. A perform_all_later counts once for the whole bulk.
  3. As each job is archived or dead-lettered, the executor deletes that execution row and bumps completed_jobs / failed_jobs. total_jobs == outstanding rows + completed + failed holds at every commit point. A job that re-enqueues itself with retry_on keeps its single row across every attempt — the batch waits for the terminal outcome, so on_success cannot fire while a retry is pending and on_failure fires if the retries exhaust.
  4. The batch finishes when no execution rows remain and the counters add up (single-winner update). A dispatcher sweep repairs crash windows — a worker that dies between archive and row-delete, an enqueue that dies between insert and send, a pending batch whose block never returned. An unsent row is only un-counted once the sweep has checked the queue and DLQ for its message; a batch is only "stalled" after config.batch_stall_threshold (default 5 minutes) without a new execution row.
  5. The dispatcher cleans up finished batches older than config.batch_retention (default 7 days).
Existing installs need rails generate pgbus:add_batch_executions (or pgbus:update). Fresh pgbus:install already includes the executions table.