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.
Installation
It's used with the Amazon SQS transport, which sends and publishes messages and creates the queues and topics.
npm i @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
npm i -D @types/aws-lambdapnpm add @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
pnpm add -D @types/aws-lambdayarn add @node-ts/bus-sqs-lambda @node-ts/bus-sqs @node-ts/bus-core
yarn add -D @types/aws-lambdaConfiguration
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.
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:
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
- Amazon SQS
- Shutting down cleanly
BusSqsLambdaReceiverin the API reference