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

ReactiveCosmosRepository删除函数失效问题求助

问题解决:Spring Boot Reactive CosmosDB 删除操作无效

核心原因

Reactive编程中所有操作都是冷流,只有当存在订阅者(调用subscribe()、block()等)时,流才会触发执行。你的代码存在以下关键问题:

  • deleteTrashcan方法直接调用deleteById但未订阅返回的Mono,导致删除操作根本未执行,日志却提前打印。
  • DataLoader中deleteAll()未与后续操作串联且未订阅,导致清空操作不执行,后续保存直接叠加在原有数据上。

修复方案

1. 修正Service中的删除方法

将返回类型改为Mono<Void>,通过链式调用处理成功/失败日志,让上层订阅执行:

@Slf4j
@Service
@AllArgsConstructor
public class TrashcanServiceImpl implements TrashcanService {
    private final TrashcanRepository trashcanRepository;
    private final TrashcanMapper trashcanMapper;

    // 其他方法保持不变...

    public Mono<Void> deleteTrashcan(String id) {
        return trashcanRepository.deleteById(id, new PartitionKey(id))
                .doOnSuccess(unused -> log.info("Deleted trashcan {}", id))
                .doOnError(error -> log.error("Failed to delete trashcan {}", id, error));
    }
}

注意:控制器层调用该方法时,需返回这个Mono,让Spring WebFlux自动处理订阅。

2. 修正DataLoader中的初始化逻辑

将deleteAll()与后续保存操作串联,确保清空完成后再执行保存,同时添加错误处理:

@Slf4j
@Component
@AllArgsConstructor
public class DataLoader {
    private final TrashcanRepository trashcanRepository;

    @PostConstruct
    void loadData() {
        Address address1 = new Address("Begijnendijk", "3130", "Liersesteenweg", "181");
        trashcanRepository.deleteAll()
                .then(trashcanRepository.save(new TrashcanDao(address1, FillStatus.EMPTY)))
                .thenMany(trashcanRepository.findAll())
                .subscribe(
                        trashcan -> log.info("Loaded trashcan with id: {}", trashcan.getId()),
                        error -> log.error("Failed to load initial data", error)
                );
    }
}

3. 优化创建方法(可选,更规范)

原创建方法直接subscribe()可能导致id未生成就返回,改为返回Mono<String>让上层处理:

public Mono<String> createTrashcan(Trashcan trashcan) {
    TrashcanDao saveTrashcan = trashcanMapper.fromTrashcanToDao(trashcan);
    return trashcanRepository.save(saveTrashcan)
            .map(TrashcanDao::getId)
            .doOnSuccess(id -> log.info("Created trashcan with id {}", id))
            .doOnError(error -> log.error("Failed to create trashcan", error));
}

关键说明

  • Reactive操作必须被订阅才会执行,不要直接调用返回Mono/Flux的方法却不处理结果。
  • 使用doOnSuccess/doOnError可在不改变流的前提下添加日志和副作用处理。
  • 异步操作需通过链式调用(then()/flatMap()等)保证执行顺序,避免并行执行导致逻辑错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 12:15:31