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:
// 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
EncodeVideocommand 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:
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.