You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将DynamoDB写入请求放入队列?求Python/Node.js Lambda示例

Got it, let's walk through how to implement this queue-based throttling for DynamoDB writes using AWS Lambda (both Python and Node.js). The core idea is exactly what you described: offload your write requests to an SQS queue first, then have a dedicated Lambda process that queue at a rate that stays within your allocated WCUs—no more blocking or retries in your original write Lambda.

Overall Architecture Breakdown

  • Producer Lambda: Takes your 2000+ data records and sends them to an SQS queue (buffering the write load)
  • SQS Queue: Acts as a buffer with configurable batch sizes and visibility timeouts to control flow
  • Consumer Lambda: Triggered automatically by SQS, pulls messages from the queue, and writes to DynamoDB at a rate aligned with your 100 WCUs

Python Implementation

1. Producer Lambda (Send Data to SQS)

This Lambda handles the initial batch of data and pushes it to SQS—no need to worry about DynamoDB limits here. We'll use boto3 and batch messages to stay within SQS limits (max 10 messages per batch, 256KB per message).

import boto3
import json

sqs = boto3.client('sqs')
QUEUE_URL = 'YOUR_SQS_QUEUE_URL'

def lambda_handler(event, context):
    # Assume event contains your 2000 data records (adjust based on your input format)
    data_records = event.get('records', [])
    
    # Split into batches of 10 (SQS max batch size)
    batches = [data_records[i:i+10] for i in range(0, len(data_records), 10)]
    
    for batch in batches:
        entries = []
        for idx, record in enumerate(batch):
            entries.append({
                'Id': str(idx),
                'MessageBody': json.dumps(record)
            })
        
        # Send batch to SQS
        sqs.send_message_batch(
            QueueUrl=QUEUE_URL,
            Entries=entries
        )
    
    return {
        'statusCode': 200,
        'body': f"Successfully queued {len(data_records)} records"
    }

2. Consumer Lambda (Process SQS & Write to DynamoDB)

This Lambda is triggered by SQS (configure the trigger in AWS Console with batch size matching your WCU capacity). We'll use DynamoDB's batch_writer to optimize writes, and ensure we stay under 100 WCUs.

For context: If your records are <1KB each, 100 WCUs allows 100 strong-consistent writes/sec or 200 eventually-consistent writes/sec. Adjust the batch size and Lambda concurrency accordingly.

import boto3
import json

dynamodb = boto3.resource('dynamodb')
TABLE_NAME = 'YOUR_DYNAMODB_TABLE_NAME'
table = dynamodb.Table(TABLE_NAME)

def lambda_handler(event, context):
    # Process each message from SQS batch
    with table.batch_writer() as batch:
        for record in event['Records']:
            item = json.loads(record['body'])
            batch.put_item(Item=item)
    
    return {
        'statusCode': 200,
        'body': f"Processed {len(event['Records'])} records"
    }

Key Configs for Consumer:

  • In SQS trigger settings: Set Batch size to 20 (if using eventually-consistent writes, that's 20 writes = 10 WCUs per batch)
  • Set Lambda Reserved concurrency to 5: 5 batches/sec × 20 records = 100 writes/sec (using eventually-consistent, that's 50 WCUs—leave room for variability)

Node.js Implementation

1. Producer Lambda (Send Data to SQS)

Using AWS SDK v3 for Node.js, we'll batch records and send them to SQS.

import { SQSClient, SendMessageBatchCommand } from "@aws-sdk/client-sqs";

const sqsClient = new SQSClient({});
const QUEUE_URL = "YOUR_SQS_QUEUE_URL";

export const handler = async (event) => {
    const dataRecords = event.records || [];
    const batches = [];
    
    // Split into batches of 10
    for (let i = 0; i < dataRecords.length; i += 10) {
        batches.push(dataRecords.slice(i, i + 10));
    }
    
    for (const batch of batches) {
        const entries = batch.map((record, idx) => ({
            Id: idx.toString(),
            MessageBody: JSON.stringify(record)
        }));
        
        const command = new SendMessageBatchCommand({
            QueueUrl: QUEUE_URL,
            Entries: entries
        });
        
        await sqsClient.send(command);
    }
    
    return {
        statusCode: 200,
        body: JSON.stringify(`Successfully queued ${dataRecords.length} records`)
    };
};

2. Consumer Lambda (Process SQS & Write to DynamoDB)

Using the DynamoDB DocumentClient to handle batch writes, aligned with your WCU limit.

import { DynamoDBClient } from "@aws-sdk/client-dynamodb";
import { DynamoDBDocumentClient, BatchWriteCommand } from "@aws-sdk/lib-dynamodb";

const dbClient = new DynamoDBClient({});
const docClient = DynamoDBDocumentClient.from(dbClient);
const TABLE_NAME = "YOUR_DYNAMODB_TABLE_NAME";

export const handler = async (event) => {
    const putRequests = event.Records.map(record => ({
        PutRequest: {
            Item: JSON.parse(record.body)
        }
    }));
    
    // Split into batches of 25 (DynamoDB batch write max)
    const batches = [];
    for (let i = 0; i < putRequests.length; i += 25) {
        batches.push(putRequests.slice(i, i + 25));
    }
    
    for (const batch of batches) {
        const command = new BatchWriteCommand({
            RequestItems: {
                [TABLE_NAME]: batch
            }
        });
        
        await docClient.send(command);
    }
    
    return {
        statusCode: 200,
        body: JSON.stringify(`Processed ${event.Records.length} records`)
    };
};

Key Configs for Consumer:

  • SQS trigger Batch size: 20 (same as Python example)
  • Lambda Reserved concurrency: 5 to cap throughput at 100 writes/sec (eventually-consistent)

Critical Notes to Avoid Issues

  • SQS Visibility Timeout: Set it to 2x your Lambda's timeout (e.g., if Lambda runs for 1min, set timeout to 2min) to prevent duplicate processing.
  • Dead-Letter Queue (DLQ): Attach a DLQ to your SQS queue to catch failed messages—this lets you debug records that couldn't be written to DynamoDB.
  • WCU Calculation: If your records are larger than 1KB, each write will consume multiple WCUs. For example, a 2KB record uses 2 WCUs for a strong-consistent write. Adjust your batch size/concurrency accordingly.
  • Eventual vs Strong Consistency: Use eventually-consistent writes (default for put_item unless specified) to double your effective WCU throughput.

内容的提问来源于stack exchange,提问作者WeCanBeFriends

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 07:13:59