Spring Boot REST API如何流式从DynamoDB获取100-200k条数据(不加载全量到内存)
从DynamoDB流式获取大量数据并返回给API消费者
要解决DynamoDB取100-200k条记录不占满内存的问题,核心是利用DynamoDB的分页查询机制结合Spring的流式输出能力,不用一次性把所有数据加载到内存里。下面是具体实现方案:
核心思路
DynamoDB的Scan或Query接口本身就支持分页(通过LastEvaluatedKey标记下一页起始位置),我们可以循环调用分页接口,每次只加载一页数据,然后立刻把这页数据输出给客户端,再加载下一页,直到所有数据处理完成。
步骤1:使用AWS SDK v2(推荐)
AWS SDK v2原生支持异步操作和流式处理,比v1更适合这种场景。先引入依赖(Maven示例):
<dependency> <groupId>software.amazon.awssdk</groupId> <artifactId>dynamodb</artifactId> </dependency> <dependency> <groupId>software.amazon.awssdk</groupId> <artifactId>netty-nio-client</artifactId> </dependency>
配置异步客户端:
@Configuration public class DynamoDbConfig { @Bean public DynamoDbAsyncClient dynamoDbAsyncClient() { return DynamoDbAsyncClient.builder() .region(Region.US_EAST_1) // 替换成你的区域 .build(); } }
步骤2:实现流式查询逻辑
场景:查询圣诞周末商务舱的乘客数据
假设你的表有bookingDate(日期)、cabinClass(舱位)字段,且建了全局二级索引(GSI):cabinClass-bookingDate-index,用cabinClass作为分区键,bookingDate作为排序键,这样用Query比Scan高效得多。
方案A:Spring WebFlux(响应式)
用Flux来发射每条数据,框架自动处理流式输出:
@Service public class PassengerService { private final DynamoDbAsyncClient dynamoDbClient; private static final String TABLE_NAME = "PassengerBookings"; private static final String GSI_NAME = "cabinClass-bookingDate-index"; public PassengerService(DynamoDbAsyncClient dynamoDbClient) { this.dynamoDbClient = dynamoDbClient; } public Flux<Passenger> getChristmasBusinessClassPassengers() { // 定义查询条件:商务舱,圣诞周末(比如2023-12-23到2023-12-25) QueryRequest initialRequest = QueryRequest.builder() .tableName(TABLE_NAME) .indexName(GSI_NAME) .keyConditionExpression("cabinClass = :cabin AND bookingDate BETWEEN :startDate AND :endDate") .expressionAttributeValues(Map.of( ":cabin", AttributeValue.builder().s("BUSINESS").build(), ":startDate", AttributeValue.builder().s("2023-12-23").build(), ":endDate", AttributeValue.builder().s("2023-12-25").build() )) .limit(1000) // 每页1000条,根据单条记录大小调整,不超过1MB .build(); // 用Flux.generate处理分页循环 return Flux.generate(() -> initialRequest, (request, sink) -> { dynamoDbClient.query(request) .thenAccept(response -> { // 把当前页的每条数据转成Passenger对象并发射 response.items().forEach(item -> { Passenger passenger = mapItemToPassenger(item); sink.next(passenger); }); // 检查是否有下一页 if (response.lastEvaluatedKey() == null) { sink.complete(); } }) .exceptionally(e -> { sink.error(e); return null; }); // 如果有下一页,构造新的查询请求 return response.lastEvaluatedKey() != null ? request.toBuilder().exclusiveStartKey(response.lastEvaluatedKey()).build() : null; }); } // 把DynamoDB的Item转成Passenger实体 private Passenger mapItemToPassenger(Map<String, AttributeValue> item) { return Passenger.builder() .id(item.get("id").s()) .name(item.get("name").s()) .bookingDate(item.get("bookingDate").s()) .cabinClass(item.get("cabinClass").s()) // 其他字段 .build(); } }
Controller层直接返回Flux:
@RestController @RequestMapping("/passengers") public class PassengerController { private final PassengerService passengerService; public PassengerController(PassengerService passengerService) { this.passengerService = passengerService; } @GetMapping("/christmas-business") public Flux<Passenger> getChristmasBusinessClassPassengers() { return passengerService.getChristmasBusinessClassPassengers(); } }
方案B:传统Spring MVC
用StreamingResponseBody实现流式输出,避免内存溢出:
@Service public class PassengerService { private final DynamoDbAsyncClient dynamoDbClient; // 同上面的配置和常量... public void streamChristmasBusinessClassPassengers(OutputStream outputStream) throws IOException { ObjectMapper objectMapper = new ObjectMapper(); QueryRequest request = QueryRequest.builder() .tableName(TABLE_NAME) .indexName(GSI_NAME) .keyConditionExpression("cabinClass = :cabin AND bookingDate BETWEEN :startDate AND :endDate") .expressionAttributeValues(Map.of( ":cabin", AttributeValue.builder().s("BUSINESS").build(), ":startDate", AttributeValue.builder().s("2023-12-23").build(), ":endDate", AttributeValue.builder().s("2023-12-25").build() )) .limit(1000) .build(); boolean hasMorePages = true; while (hasMorePages) { // 同步查询(如果用异步可以调整,但同步更简单适合MVC场景) QueryResponse response = dynamoDbClient.query(request).join(); // 把当前页的每条数据写入输出流 for (Map<String, AttributeValue> item : response.items()) { Passenger passenger = mapItemToPassenger(item); objectMapper.writeValue(outputStream, passenger); outputStream.write('\n'); // 每条数据换行,方便客户端解析 outputStream.flush(); // 及时刷出,避免缓存 } // 检查下一页 if (response.lastEvaluatedKey() == null) { hasMorePages = false; } else { request = request.toBuilder().exclusiveStartKey(response.lastEvaluatedKey()).build(); } } } }
Controller层返回StreamingResponseBody:
@RestController @RequestMapping("/passengers") public class PassengerController { private final PassengerService passengerService; public PassengerController(PassengerService passengerService) { this.passengerService = passengerService; } @GetMapping(value = "/christmas-business", produces = MediaType.APPLICATION_NDJSON_VALUE) public StreamingResponseBody streamPassengers() { return outputStream -> { passengerService.streamChristmasBusinessClassPassengers(outputStream); }; } }
这里用APPLICATION_NDJSON_VALUE(换行分隔JSON),客户端可以逐行解析,不用等全量数据。
关键注意事项
- 索引优化:一定要用
Query+GSI代替Scan,Scan会全表扫描,性能极差,尤其是大数据量场景。针对你的查询条件,把常用过滤字段设为GSI的分区/排序键。 - 分页大小:DynamoDB每页最多返回1MB数据,所以
limit值要根据单条记录大小调整,比如单条1KB的话,设1000就刚好1MB,避免频繁请求。 - 内存控制:每次只处理一页数据,处理完立刻输出,不要把所有页的数据存到集合里。
- 异常处理:分页过程中如果出现DynamoDB限流(ProvisionedThroughputExceededException),要加重试逻辑(可以用SDK自带的重试器)。
- 客户端兼容性:如果返回NDJSON,要提前和客户端沟通格式;如果客户端需要数组格式,MVC场景可以先输出
[,然后每条数据加逗号(最后一条不加),最后输出],但这种方式要注意异常情况下的格式完整性。
内容的提问来源于stack exchange,提问作者xyzmar
相关产品推荐
相关产品推荐

