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

如何基于主键实现DynamoDB记录归档?含Lambda脚本改造需求

解决方案:DynamoDB主键冲突时归档旧数据到S3并写入新数据

要实现你的需求,需要修改原Lambda脚本,在写入新数据前先检查DynamoDB中是否存在对应主键的旧记录,若存在则归档至S3,再执行新数据写入。以下是具体实现步骤和修改后的代码:

核心修改点

  • 移除批量写入的batch_writer,改为逐个处理每条记录,以便单独查询旧数据
  • 新增S3客户端,用于归档旧数据
  • 对每条Oracle同步来的记录,先查询DynamoDB中是否存在对应主键的记录
  • 若存在旧记录,将其序列化为JSON格式后上传至指定S3桶
  • 最后将新数据写入DynamoDB

修改后的完整代码

from boto3 import resource, client
from cx_Oracle import connect
from logging import getLogger, DEBUG
from datetime import datetime, timedelta
import json
import uuid

def lambda_handler(event, context):
    log_g_name = "/aws/lambda/lambda"
    client_log = client('logs')
    
    # 获取上次Lambda执行时间(加5秒避免遗漏)
    stream_response = client_log.describe_log_streams(
        logGroupName=log_g_name,
        orderBy='LastEventTime',
        limit=1,
        descending=True
    )
    last_run = datetime.fromtimestamp(
        stream_response['logStreams'][0]['lastEventTimestamp']/1000.0
    ) + timedelta(seconds=5)
    
    # 初始化DynamoDB和S3客户端
    dynamodb = resource('dynamodb')
    table = dynamodb.Table('table_name')
    s3 = client('s3')
    s3_bucket_name = 'your-archive-bucket-name'  # 替换为你的S3归档桶名
    
    # 连接Oracle并获取增量数据
    connstr = 'conn_str'
    connection = connect(connstr)
    cur = connection.cursor()
    cur.execute("select * from table where timestamp >= '%s'" % last_run)
    
    columns = [column[0] for column in cur.description]
    results = []
    for row in cur.fetchall():
        results.append(dict(zip(columns, row)))
    
    # 处理每条记录:归档旧数据+写入新数据
    for item in results:
        # 替换为你DynamoDB表实际的主键字段(示例为'id')
        primary_key = {'id': item['id']}
        
        # 查询DynamoDB中的旧记录
        response = table.get_item(Key=primary_key)
        if 'Item' in response:
            old_item = response['Item']
            
            # 处理datetime类型序列化(转为ISO格式字符串)
            for key, value in old_item.items():
                if isinstance(value, datetime):
                    old_item[key] = value.isoformat()
            
            # 生成唯一S3对象键,避免归档文件重复覆盖
            s3_object_key = f"dynamodb-archive/{item['id']}_{datetime.now().strftime('%Y%m%d%H%M%S')}_{uuid.uuid4().hex}.json"
            
            # 将旧数据上传至S3
            s3.put_object(
                Bucket=s3_bucket_name,
                Key=s3_object_key,
                Body=json.dumps(old_item),
                ContentType='application/json'
            )
        
        # 写入新数据到DynamoDB
        table.put_item(Item=item)
    
    # 关闭数据库连接
    cur.close()
    connection.close()

注意事项

  • 主键字段替换:代码中假设主键为id,请替换为你DynamoDB表实际的主键字段名
  • S3桶配置:将your-archive-bucket-name替换为你的归档S3桶名称
  • 权限配置:确保Lambda角色拥有以下权限:
    • DynamoDB的GetItem和PutItem权限
    • S3的PutObject权限
    • CloudWatch Logs的DescribeLogStreams权限
  • 序列化处理:如果Oracle返回的字段包含Decimal等特殊类型,需要额外处理序列化逻辑,避免JSON转换报错
  • 性能优化:若同步数据量较大,可考虑结合batch_get_item批量查询旧数据,再批量归档后批量写入,平衡性能与代码复杂度

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 02:50:02