如何用Webflux获取DynamoDB全局二级索引的全部数据?
解决Webflux查询DynamoDB全局二级索引无法获取全部数据的问题
问题原因
DynamoDB的query操作单次返回结果受限于1MB的数据大小限制(或最多1000条记录,以先触发的条件为准),当结果集超过这个阈值时,会返回lastEvaluatedKey标记下一页的起始位置。你的代码仅执行了一次query调用,只获取了第一页的1700条数据,没有处理后续分页数据。
解决方案
利用Webflux的Flux.expand操作符实现异步递归分页查询,直到获取所有数据为止。修改后的代码如下:
var queryConditional = QueryConditional.keyEqualTo(Key.builder().partitionValue(status).build()); var index = dynamoDbAsyncTable.index(Attributes.STATUS); // 构建初始查询请求 QueryEnhancedRequest initialRequest = QueryEnhancedRequest.builder() .queryConditional(queryConditional) .build(); // 递归获取所有分页数据 Flux<Table> allItems = Flux.from(index.query(initialRequest)) .expand(page -> { // 检查是否存在下一页的起始标记 if (page.lastEvaluatedKey() != null) { QueryEnhancedRequest nextRequest = QueryEnhancedRequest.builder() .queryConditional(queryConditional) .exclusiveStartKey(page.lastEvaluatedKey()) .build(); return Flux.from(index.query(nextRequest)); } // 无更多数据时终止递归 return Flux.empty(); }) // 将每页的元素列表展开为单个流元素 .flatMapIterable(Page::items); // 收集所有元素并转换为目标集合 return allItems.collect(Collectors.toSet()) .map(courierTables -> courierTables.stream() .map(Table::toCourier) .collect(Collectors.toSet()));
关键说明
- 使用
Flux替代Mono:因为需要处理多页数据的流式返回,而非单页结果 expand操作符:自动递归处理分页逻辑,每次基于上一页的lastEvaluatedKey构建下一页请求flatMapIterable:将每个Page中的元素列表展开为单个元素的流,方便后续统一收集- 内存提示:如果数据量极大(如20000条),一次性加载到内存可能引发OOM,建议根据业务场景考虑分批次返回给前端
内容的提问来源于stack exchange,提问作者Anjali Verma
相关产品推荐
相关产品推荐

