Skip to content

Long running processes ​

Some tasks can't finish within the time a message can be held, such as encoding a video or preparing a large file for download. This page explains why they shouldn't run in a handler, and an approach that works.

Why not in a handler ​

When a message is read from the queue, it has to be handled and deleted within a timeout, usually 30 seconds to a few minutes. If it isn't, the queue assumes the consumer died, and makes the message visible again for another consumer, which starts the same work. After a number of attempts, the message goes to the dead letter queue.

A long task in a handler also blocks that worker from handling other messages while it waits. Handlers should finish as quickly as they can.

Say a command, EncodeVideo, can take up to an hour. Its handler can't wait for the encoding to finish. Instead, it should start the work somewhere else, and events should report its progress.

A naive approach ​

One way that's not recommended is to run the task in the background:

ts
// Not recommended: the work is lost if it fails or the service restarts
export const encodeVideoInBackgroundHandler = handlerFor(
  EncodeVideo,
  command => {
    setTimeout(() => {
      videoService.encode(command.videoId).catch(console.error)
    }, 0)
  }
)

The command is deleted straight away, but:

  • Nothing retries the task if it fails, and nothing publishes a message to say it did.
  • Nothing balances the load. One instance may receive every EncodeVideo command and run hundreds of tasks until it crashes.
  • When the service restarts, the tasks running in the background are lost and not retried.

A task per job ​

If your application runs on Kubernetes, ECS, Docker Swarm or similar, start a container task for each job, and leave it to the scheduler to place it. Scaling then follows the number of jobs. The handler starts the task and returns:

ts
export const encodeVideoHandler = handlerFor(
  EncodeVideo,
  async (command, _attributes, ctx) => {
    // Start a container task to do the work, and return straight away
    const taskId = await taskScheduler.runTask('video-encoder', [
      command.videoId
    ])
    await ctx.publish(new VideoEncodingStarted(command.videoId, taskId))
  }
)

The task publishes an event such as VideoEncoded when it's done.

Recovering from failed tasks ​

Starting tasks from handlers doesn't help when a task fails or the scheduler stops it. A workflow can track each job: start it on VideoEncodingStarted, listen for the scheduler's system messages that report a task exited, start the task again when needed, and complete on VideoEncoded.

See also ​

Released under the MIT License.