Process DynamoDB Streams in Lambda With TypeScript

Glowing blue fibre optic strands carrying light against a dark background

Photo by Mimicry Hu on Unsplash

To process DynamoDB Streams in Lambda with TypeScript, enable the stream with UpdateTable (StreamViewType: "NEW_AND_OLD_IMAGES"), connect it with CreateEventSourceMapping, and type the handler as (event: DynamoDBStreamEvent) => Promise<DynamoDBBatchResponse> from @types/aws-lambda. Convert each image with unmarshall from @aws-sdk/util-dynamodb, and return the first failed record’s SequenceNumber in batchItemFailures.

DynamoDB Streams turns every insert, update and delete in a table into an ordered change record, and Lambda can consume those records without you running any pollers. This guide is for Node.js and TypeScript developers building audit trails, search index sync, notifications or cache invalidation on top of a table.

You’ll get the three pieces as tested code: turning the stream on, an event source mapping with retry settings that won’t block a shard for a day, and a typed handler with an idempotent write and partial batch responses, plus a test that runs it against a sample stream event. Samples were type-checked with strict tsc against AWS SDK v3 3.1141.0 and @types/aws-lambda 8.10.163, and run in September 2026 with aws-sdk-client-mock.

How does Lambda read a DynamoDB stream?

  • Polling and batching. The Lambda Developer Guide says Lambda polls each shard 4 times a second and invokes your function with a batch of up to BatchSize records (default 100, maximum 10,000), within a 6 MB payload.
  • Ordering. DynamoDB Streams guarantees each record appears once, in the order the item was modified. Order is per item (its full primary key), not across the table. ParallelizationFactor (1 to 10) processes a shard with more concurrent batches and still keeps per-item order.
  • At least once. Lambda’s event source mapping can deliver a record more than once, so the handler must be idempotent.
  • 24-hour retention. Stream records are kept for 24 hours. With default settings, a record that keeps failing blocks its shard until it expires.
  • Two readers per shard. For a single-Region table, design for at most two Lambda functions on the same stream; one for global tables.

Prerequisites

  • Node.js 18 or later, a project with "type": "module", and @aws-sdk/client-dynamodb, @aws-sdk/client-lambda and @aws-sdk/util-dynamodb.
  • Dev dependencies: typescript, tsx, aws-sdk-client-mock and @types/aws-lambda, the community type definitions for Lambda events.
  • A deployed function (here orders-stream-handler) and a standard SQS queue for failed batches. Sending and receiving SQS messages with SDK v3 helps when you process that queue later.

How to process DynamoDB Streams in Lambda with TypeScript, step by step

  1. Enable the streamPick the view type up front. You can’t edit it later without turning the stream off and on, which creates a new stream ARN.
  2. Grant the execution role stream accessRead actions on the stream, plus sqs:SendMessage for the failure queue.
  3. Create the event source mappingStart at TRIM_HORIZON, cap retries and record age, bisect on error, and turn on ReportBatchItemFailures.
  4. Write the handlerType the event, unmarshall images, make writes idempotent, and stop at the first failure.
  5. Test it locallyFeed it a sample event with a mocked SDK client before deploying.

Step 1: enable the stream on the table

enable-stream.ts

// enable-stream.ts: turn on a DynamoDB stream with old and new item images and print its ARN.
// Usage: TABLE=orders npx tsx enable-stream.ts
import { DynamoDBClient, DescribeTableCommand, UpdateTableCommand, waitUntilTableExists } from "@aws-sdk/client-dynamodb";

const client = new DynamoDBClient({ region: process.env.AWS_REGION ?? "us-east-1" });
const TableName = process.env.TABLE ?? "orders";

const { Table } = await client.send(new DescribeTableCommand({ TableName }));
if (Table?.StreamSpecification?.StreamEnabled) {
  // The view type can't be edited: disable the stream and enable a new one to change it.
  console.log(`Stream already on (${Table.StreamSpecification.StreamViewType}): ${Table.LatestStreamArn}`);
} else {
  await client.send(new UpdateTableCommand({
    TableName,
    StreamSpecification: { StreamEnabled: true, StreamViewType: "NEW_AND_OLD_IMAGES" },
  }));
  await waitUntilTableExists({ client, maxWaitTime: 300 }, { TableName }); // waits for ACTIVE
  const after = await client.send(new DescribeTableCommand({ TableName }));
  console.log(`Stream enabled: ${after.Table?.LatestStreamArn}`);
}
Output

Stream enabled: arn:aws:dynamodb:us-east-1:123456789012:table/orders/stream/2027-07-23T09:14:02.118

The DynamoDB Developer Guide lists four view types: KEYS_ONLY, NEW_IMAGE, OLD_IMAGE and NEW_AND_OLD_IMAGES. Choose the last when the handler needs to compare before and after, as an audit trail does. Enabling a stream on a table that already has one returns a ValidationException, which is why the script checks first. A PutItem or UpdateItem that changes nothing writes no stream record; the guide to updating DynamoDB items with SDK v3 covers the writes that do.

Step 2: create the event source mapping

create-mapping.ts

// create-mapping.ts: connect a DynamoDB stream to a Lambda function with retries, bisecting,
// partial batch responses, a filter and an SQS on-failure destination.
// Usage: STREAM_ARN=arn:aws:dynamodb:...:table/orders/stream/... DLQ_ARN=arn:aws:sqs:...:orders-stream-dlq \
//        npx tsx create-mapping.ts
import { LambdaClient, CreateEventSourceMappingCommand } from "@aws-sdk/client-lambda";

const lambda = new LambdaClient({ region: process.env.AWS_REGION ?? "us-east-1" });
const streamArn = process.env.STREAM_ARN;
const dlqArn = process.env.DLQ_ARN;
if (!streamArn || !dlqArn) throw new Error("Set STREAM_ARN and DLQ_ARN");

// Only records whose new image is an order. REMOVE events have no NewImage, so they're filtered out too.
const orderFilter = { dynamodb: { NewImage: { entityType: { S: ["ORDER"] } } } };

const out = await lambda.send(new CreateEventSourceMappingCommand({
  FunctionName: "orders-stream-handler",
  EventSourceArn: streamArn,
  StartingPosition: "TRIM_HORIZON", // LATEST can miss records while the mapping is being created
  BatchSize: 100, // default 100, maximum 10,000
  MaximumBatchingWindowInSeconds: 5, // wait up to 5 s to fill a batch
  ParallelizationFactor: 1, // 1-10 concurrent batches per shard; per-item order is kept either way
  BisectBatchOnFunctionError: true, // split a failing batch in two before retrying
  MaximumRetryAttempts: 5, // default -1 retries until the record expires (24 hours)
  MaximumRecordAgeInSeconds: 3600, // give up on records older than an hour
  FunctionResponseTypes: ["ReportBatchItemFailures"],
  FilterCriteria: { Filters: [{ Pattern: JSON.stringify(orderFilter) }] },
  DestinationConfig: { OnFailure: { Destination: dlqArn } },
}));

console.log(`Mapping ${out.UUID} is ${out.State}`);
Output

Mapping a1b2c3d4-5678-90ab-cdef-11111EXAMPLE is Creating
Setting Default Why change it
StartingPosition required The Lambda guide warns LATEST can miss events while the mapping is created or updated; TRIM_HORIZON doesn’t.
MaximumRetryAttempts -1 (until the record expires) A bad record would otherwise block the shard for up to a day.
MaximumRecordAgeInSeconds -1 Skip records too old to matter; -1 or 60 to 604,800.
BisectBatchOnFunctionError false Splits a failing batch to isolate the bad record; splits don’t use retry attempts.
FunctionResponseTypes none ReportBatchItemFailures lets the handler say where to resume. Without it Lambda ignores batchItemFailures.
DestinationConfig.OnFailure none Where Lambda sends details of batches it gives up on.
FilterCriteria none Skip records before invoking; DynamoDB filters match only the dynamodb key.

Filter patterns compare typed values, so { "S": ["ORDER"] }, not "ORDER". The Lambda guide notes that numeric filter operators don’t work for DynamoDB, because numbers arrive as strings. The script to find Lambda SQS triggers without partial batch response audits the same setting on queue triggers.

Step 3: write a typed handler

handler.ts

// handler.ts: typed DynamoDB Streams handler that writes an audit row per order change and reports
// the first failed record so Lambda retries from there.
import type { AttributeValue as LambdaAttributeValue, DynamoDBBatchResponse, DynamoDBRecord, DynamoDBStreamEvent } from "aws-lambda";
import { ConditionalCheckFailedException, DynamoDBClient, PutItemCommand, type AttributeValue } from "@aws-sdk/client-dynamodb";
import { marshall, unmarshall } from "@aws-sdk/util-dynamodb";

export interface Order {
  pk: string;
  entityType: "ORDER";
  status: string;
  total: number;
}

const ddb = new DynamoDBClient({});
const AUDIT_TABLE = process.env.AUDIT_TABLE ?? "orders-audit";

// aws-lambda's AttributeValue and the SDK's AttributeValue describe the same JSON but are different types.
const toOrder = (image: Record<string, LambdaAttributeValue> | undefined): Order | undefined =>
  image ? (unmarshall(image as Record<string, AttributeValue>) as Order) : undefined;

async function processRecord(record: DynamoDBRecord): Promise<void> {
  const before = toOrder(record.dynamodb?.OldImage);
  const after = toOrder(record.dynamodb?.NewImage);
  if (!after || before?.status === after.status) return; // nothing to audit

  try {
    await ddb.send(new PutItemCommand({
      TableName: AUDIT_TABLE,
      Item: marshall({
        pk: after.pk,
        sk: record.eventID ?? "", // eventID makes the write idempotent across retries
        from: before?.status ?? null,
        to: after.status,
        at: record.dynamodb?.ApproximateCreationDateTime ?? 0,
      }),
      ConditionExpression: "attribute_not_exists(sk)",
    }));
  } catch (err) {
    if (err instanceof ConditionalCheckFailedException) return; // already written by an earlier attempt
    throw err;
  }
}

export const handler = async (event: DynamoDBStreamEvent): Promise<DynamoDBBatchResponse> => {
  for (const record of event.Records) {
    try {
      await processRecord(record);
    } catch (err) {
      console.error(`record ${record.eventID} failed`, err);
      // Records in a shard are ordered: stop here and let Lambda retry from this sequence number.
      return { batchItemFailures: [{ itemIdentifier: record.dynamodb?.SequenceNumber ?? "" }] };
    }
  }
  return { batchItemFailures: [] };
};

Three details carry the weight:

  • Two AttributeValue types. @types/aws-lambda and @aws-sdk/client-dynamodb each define one, and strict tsc won’t pass a stream image straight to unmarshall. The cast is safe: both describe the same wire format. Numbers come back as JavaScript numbers; pass { wrapNumbers: true } to unmarshall if they can exceed Number.MAX_SAFE_INTEGER.
  • Idempotent writes. The audit row’s sort key is the stream eventID, and attribute_not_exists(sk) turns a duplicate delivery into a harmless ConditionalCheckFailedException.
  • Stop at the first failure. The Lambda guide says it checkpoints at the lowest reported sequence number and retries every record after it, so processing later records only repeats work. An empty batchItemFailures list means success; an empty or null itemIdentifier fails the whole batch.

If the handler must update several items atomically, DynamoDB TransactWriteItems with SDK v3 covers transactions, including how they show up in a stream.

Step 4: test the handler with a sample stream event

handler.test.ts

// handler.test.ts: run the handler against a sample stream event with a mocked DynamoDB client.
// Usage: npx tsx handler.test.ts
import { mockClient } from "aws-sdk-client-mock";
import { ConditionalCheckFailedException, DynamoDBClient, PutItemCommand } from "@aws-sdk/client-dynamodb";
import type { DynamoDBRecord, DynamoDBStreamEvent } from "aws-lambda";
import { handler } from "./handler.js";

const ddb = mockClient(DynamoDBClient);
const streamArn = "arn:aws:dynamodb:us-east-1:123456789012:table/orders/stream/2027-07-23T09:14:02.118";

const order = (status: string) => ({
  pk: { S: "ORDER#1001" }, entityType: { S: "ORDER" }, status: { S: status }, total: { N: "42.5" },
});
const record = (eventID: string, seq: string, from: string | null, to: string): DynamoDBRecord => ({
  eventID,
  eventName: from ? "MODIFY" : "INSERT",
  eventSource: "aws:dynamodb",
  eventSourceARN: streamArn,
  awsRegion: "us-east-1",
  dynamodb: {
    ApproximateCreationDateTime: 1784797200,
    Keys: { pk: { S: "ORDER#1001" } },
    ...(from ? { OldImage: order(from) } : {}),
    NewImage: order(to),
    SequenceNumber: seq,
    SizeBytes: 120,
    StreamViewType: "NEW_AND_OLD_IMAGES",
  },
});

const event: DynamoDBStreamEvent = {
  Records: [
    record("e1", "100000000000000000001", null, "PLACED"),
    record("e2", "100000000000000000002", "PLACED", "PLACED"), // no status change: skipped
    record("e3", "100000000000000000003", "PLACED", "PAID"),
    record("e4", "100000000000000000004", "PAID", "SHIPPED"),
  ],
};

// 1. A retry where e3's audit row already exists: the conditional failure counts as success.
ddb.on(PutItemCommand)
  .resolvesOnce({})
  .rejectsOnce(new ConditionalCheckFailedException({ message: "mock: row exists", $metadata: {} }))
  .resolves({});
console.log("all ok:", JSON.stringify(await handler(event)));
console.log("PutItem calls:", ddb.commandCalls(PutItemCommand).length);

// 2. The write for e4 fails: the handler reports e4's sequence number.
ddb.reset();
ddb.on(PutItemCommand).resolvesOnce({}).resolvesOnce({}).rejectsOnce(new Error("mock: write failed"));
console.log("e4 fails:", JSON.stringify(await handler(event)));
Output

all ok: {"batchItemFailures":[]}
PutItem calls: 3
record e4 failed Error: mock: write failed
e4 fails: {"batchItemFailures":[{"itemIdentifier":"100000000000000000004"}]}

The stack trace from console.error is trimmed. e2 doesn’t change the status, so only three writes happen; the duplicate e3 counts as success; and when e4‘s write fails, the handler returns its sequence number so Lambda retries from there. The guide to mocking AWS SDK v3 in unit tests shows the same technique in Jest and Vitest.

What happens when a batch keeps failing?

Lambda retries until the batch succeeds, the records exceed MaximumRecordAgeInSeconds, or retries reach MaximumRetryAttempts. Then it discards the batch and sends an invocation record to the on-failure destination. For SQS and SNS destinations, that record holds metadata only: DDBStreamBatchInfo with the shard ID, start and end sequence numbers and stream ARN, not the items. You have to read them back from the stream before they expire after 24 hours; an S3 destination also includes the original payload. This differs from asynchronous invocations, which the script to find Lambda functions without a failure destination covers. When the function itself errors, investigating Lambda errors with CloudWatch is the next step.

Cost is mostly Lambda’s: the DynamoDB pricing page says GetRecords calls made by Lambda triggers aren’t charged, unless the function runs on Lambda Managed Instances. Other stream readers pay per read request unit; the AWS Price List for DynamoDB (published 11 September 2026) shows 2.5 million free a month in US East (N. Virginia), then $0.0000002 each.

Which IAM permissions are needed?

The function’s execution role reads the stream, writes the audit table and sends failures to the queue. The actions match the AWS managed policy AWSLambdaDynamoDBExecutionRole, scoped down:

orders-stream-handler-role-policy.json

{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Sid": "ReadOrdersStream",
      "Effect": "Allow",
      "Action": [
        "dynamodb:DescribeStream",
        "dynamodb:GetRecords",
        "dynamodb:GetShardIterator"
      ],
      "Resource": "arn:aws:dynamodb:us-east-1:123456789012:table/orders/stream/*"
    },
    {
      "Sid": "ListStreams",
      "Effect": "Allow",
      "Action": "dynamodb:ListStreams",
      "Resource": "*"
    },
    {
      "Sid": "WriteAudit",
      "Effect": "Allow",
      "Action": "dynamodb:PutItem",
      "Resource": "arn:aws:dynamodb:us-east-1:123456789012:table/orders-audit"
    },
    {
      "Sid": "SendFailedBatches",
      "Effect": "Allow",
      "Action": "sqs:SendMessage",
      "Resource": "arn:aws:sqs:us-east-1:123456789012:orders-stream-dlq"
    },
    {
      "Sid": "Logs",
      "Effect": "Allow",
      "Action": ["logs:CreateLogGroup", "logs:CreateLogStream", "logs:PutLogEvents"],
      "Resource": "arn:aws:logs:us-east-1:123456789012:log-group:/aws/lambda/orders-stream-handler*"
    }
  ]
}

Whoever runs the setup scripts needs dynamodb:DescribeTable, dynamodb:UpdateTable and lambda:CreateEventSourceMapping. The IAM policy generator for TypeScript code lists the actions a handler uses; if the table holds data you can’t lose, enabling DynamoDB point-in-time recovery is worth doing too.

Troubleshooting and limits

  • IteratorAge keeps growing. The function is too slow or failing. Raise ParallelizationFactor, fix the failing record, or lower MaximumRecordAgeInSeconds.
  • The handler returns failures but everything is retried. ReportBatchItemFailures isn’t enabled on the mapping, or itemIdentifier is empty.
  • Records are missing after deploy. The mapping started at LATEST, or a filter dropped them; filtered-out records aren’t sent to the function.
  • Throttling on the stream. More than two consumers read the same shard. Fan out from one function instead, for example to EventBridge or SQS.
  • Large bundle, slow cold start. Import only the clients you use; reducing AWS SDK v3 Lambda bundle size shows how.
  • Cross-Region triggers. Not supported: the function and stream must be in the same Region.

Frequently asked questions

What type is the DynamoDB stream event in TypeScript?

DynamoDBStreamEvent from @types/aws-lambda, with DynamoDBRecord for each record and DynamoDBBatchResponse as the return type when partial batch responses are on.

How long are DynamoDB stream records kept?

24 hours. Records older than that can be trimmed at any time, so failed batches must be handled within a day.

Does Lambda process DynamoDB stream records in order?

Yes, per item. Changes to the same primary key arrive in order, including with a parallelization factor above 1. There’s no ordering across different items.

How do I convert NewImage to a plain object?

Use unmarshall from @aws-sdk/util-dynamodb, casting the image to the SDK’s AttributeValue record type because @types/aws-lambda declares its own.

Can two Lambda functions read the same DynamoDB stream?

Yes, up to two per shard for a single-Region table. For global tables, AWS recommends one.

Related guides

Ask your AWS account in plain English

Your first 15 runs are free, with no OpenAI key needed.

npx chatwithcloud