Skip to content

SQS and Lambda ​

@node-ts/bus-sqs-lambda lets an AWS Lambda function that's triggered by an SQS queue pass each batch of messages to the bus, instead of the bus polling the queue itself. Handlers, workflows and retries work as they do in a long running service. This page covers configuring it and handling partial batch failures.

@node-ts/bus-sqs-lambda@node-ts/bus-sqs-lambda version on npmSource

Installation ​

It's used with the Amazon SQS transport, which sends and publishes messages and creates the queues and topics.

sh
npm i @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
npm i -D @types/aws-lambda
sh
pnpm add @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
pnpm add -D @types/aws-lambda
sh
yarn add @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
yarn add -D @types/aws-lambda

Configuration ​

Configure the bus with the SQS transport and a BusSqsLambdaReceiver, initialize it when the module loads, and pass each event to bus.receive(). Don't call bus.start(): Lambda reads the queue instead.

ts
const sqsTransport = new SqsTransport({
  awsRegion: process.env.AWS_REGION,
  awsAccountId: process.env.AWS_ACCOUNT_ID,
  queueName: 'reservations-service',
  deadLetterQueueName: 'reservations-service-dead-letter'
})

const bus = Bus.configure()
  .withMessageTypes(messageTypes)
  .withTransport(sqsTransport)
  .withHandler(reserveRoomHandler)
  // Lambda reads the queue and passes each batch to the bus
  .withReceiver(new BusSqsLambdaReceiver())
  // Lambda owns the process, so don't listen for shutdown signals
  .withInterruptSignals([])
  .build()

// Runs once per Lambda instance, when the module is loaded
await bus.initialize()

// Pass a function, rather than bus.receive itself, so it keeps its `this`
export const handler: SQSHandler = event => bus.receive(event)

Each record is dispatched to its handlers, at most withConcurrency() at a time. Records that succeed are left for Lambda to delete. Records whose message has no handler are discarded, not retried. A record whose handler calls ctx.returnMessage() is treated as failed, so that Lambda retries it.

By default, if any record fails, bus.receive() rejects once the whole batch has been handled, and Lambda retries the whole batch, including the records that succeeded.

Partial batch failures ​

To retry only the records that failed, pass reportBatchItemFailures: true. bus.receive() then resolves with an SQSBatchResponse listing the failed records instead of rejecting:

ts
const batchBus = Bus.configure()
  .withMessageTypes(messageTypes)
  .withTransport(
    new SqsTransport({
      awsRegion: process.env.AWS_REGION,
      awsAccountId: process.env.AWS_ACCOUNT_ID,
      queueName: 'reservations-service'
    })
  )
  .withHandler(reserveRoomHandler)
  .withReceiver(new BusSqsLambdaReceiver({ reportBatchItemFailures: true }))
  .withInterruptSignals([])
  .build()

await batchBus.initialize()

// Resolves with the records that failed, so Lambda only retries those
export const batchHandler: SQSHandler = event =>
  batchBus.receive<SQSBatchResponse>(event)

WARNING

The Lambda's SQS event source mapping must include ReportBatchItemFailures in its FunctionResponseTypes. Without it, Lambda ignores the response, and deletes the failed records along with the rest of the batch.

Records are handled concurrently, so on a FIFO queue a failed record doesn't stop later records in the same message group from being handled.

See also ​

Released under the MIT License.