如何基于主键实现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权限
- DynamoDB的
- 序列化处理:如果Oracle返回的字段包含Decimal等特殊类型,需要额外处理序列化逻辑,避免JSON转换报错
- 性能优化:若同步数据量较大,可考虑结合
batch_get_item批量查询旧数据,再批量归档后批量写入,平衡性能与代码复杂度
内容的提问来源于stack exchange,提问作者user20986502
相关产品推荐
相关产品推荐

