如何将AWS DynamoDB数据同步至MongoDB Compass并实时更新?
解决方案:DynamoDB数据导入MongoDB Compass + 实时同步
一、一次性将DynamoDB数据导入MongoDB Compass
1. 导出DynamoDB数据
使用AWS CLI命令扫描目标表并导出为JSON文件:
aws dynamodb scan --table-name 你的表名 --output json > dynamodb-data.json
数据量较大时,可添加
--page-size和--max-items参数分批导出,或使用DynamoDB批量导出到S3后再下载。
2. 转换DynamoDB数据格式
DynamoDB导出的JSON包含类型标识(如{"S": "字符串值"}),需转为MongoDB兼容的普通JSON。示例Python转换脚本:
import json with open('dynamodb-data.json', 'r') as f: data = json.load(f) converted_items = [] for item in data['Items']: converted_item = {} for key, value in item.items(): if 'S' in value: converted_item[key] = value['S'] elif 'N' in value: converted_item[key] = float(value['N']) if '.' in value['N'] else int(value['N']) elif 'BOOL' in value: converted_item[key] = value['BOOL'] # 可根据需求扩展其他类型(如L、M等) converted_items.append(converted_item) with open('mongodb-data.json', 'w') as f: json.dump(converted_items, f, indent=2)
3. 导入到MongoDB Compass
- 打开MongoDB Compass,连接本地MongoDB实例。
- 选择目标数据库和集合(不存在则创建)。
- 点击顶部导入数据按钮,选择转换后的
mongodb-data.json文件。 - 按向导配置字段匹配规则,完成导入。
二、实现DynamoDB数据变更实时同步到MongoDB
方案1:基于DynamoDB Streams(推荐,支持全场景同步)
适合所有修改DynamoDB数据的场景(含外部系统操作)。
开启DynamoDB表流功能
- 登录AWS控制台,进入目标DynamoDB表详情页,切换到流标签。
- 点击启用,选择捕获类型(如“新旧图像”,用于完整获取变更前后数据),保存配置。
SpringBoot集成Stream监听器
- 添加Maven依赖:
<dependency> <groupId>org.springframework.cloud</groupId> <artifactId>spring-cloud-starter-aws-messaging</artifactId> </dependency> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-dynamodb</artifactId> </dependency> - 在
application.yml配置AWS凭证与区域:cloud: aws: credentials: access-key: 你的AWS访问密钥 secret-key: 你的AWS秘密密钥 region: static: 你的DynamoDB区域(如us-east-1) - 创建监听器类处理变更事件:
import com.amazonaws.services.dynamodbv2.model.Record; import com.amazonaws.services.dynamodbv2.model.StreamRecord; import org.springframework.cloud.aws.messaging.listener.annotation.DynamoDBStreamListener; import org.springframework.stereotype.Component; import java.util.Map; @Component public class DynamoDBStreamSyncListener { private final YourMongoRepository mongoRepository; public DynamoDBStreamSyncListener(YourMongoRepository mongoRepository) { this.mongoRepository = mongoRepository; } @DynamoDBStreamListener("你的DynamoDB表流ARN") public void handleStreamEvent(Record record) { StreamRecord streamRecord = record.getDynamodb(); Map<String, Object> newImage = streamRecord.getNewImage(); String eventName = record.getEventName(); YourMongoEntity entity = convertToMongoEntity(newImage); switch (eventName) { case "INSERT": case "MODIFY": mongoRepository.save(entity); break; case "REMOVE": mongoRepository.deleteById(entity.getId()); break; } } // 实现DynamoDB数据到MongoDB实体的转换逻辑 private YourMongoEntity convertToMongoEntity(Map<String, Object> dynamoDbItem) { YourMongoEntity entity = new YourMongoEntity(); entity.setName(((Map<String, String>) dynamoDbItem.get("name")).get("S")); // 扩展其他字段转换... return entity; } } - 注意事项:添加重试机制(如Spring Retry)处理同步失败;实现幂等性避免重复处理;大流量场景下可异步处理。
- 添加Maven依赖:
方案2:业务层同步(仅适用于当前SpringBoot应用为唯一数据写入源)
如果只有当前应用操作DynamoDB数据,可在业务层直接同步:
@Service public class DataSyncService { private final YourDynamoDbRepository dynamoDbRepository; private final YourMongoRepository mongoRepository; public DataSyncService(YourDynamoDbRepository dynamoDbRepository, YourMongoRepository mongoRepository) { this.dynamoDbRepository = dynamoDbRepository; this.mongoRepository = mongoRepository; } public void saveData(YourDataEntity data) { dynamoDbRepository.save(data); YourMongoEntity mongoEntity = convertToMongoEntity(data); mongoRepository.save(mongoEntity); } // 同理实现更新、删除方法... }
内容的提问来源于stack exchange,提问作者Andreea Nita
相关产品推荐
相关产品推荐

