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

Azure Cosmos DB Java SDK:按日期范围/续令牌重处理变更源

问题描述

我正在使用Azure Cosmos DB NoSQL API,需要借助Java SDK在指定日期范围或两个续令牌(continuation tokens)之间重处理变更源(Change Feed)中的变更。
查阅文档后,我未找到按日期范围筛选变更的直接选项。我的业务场景需要获取特定时间戳之间的变更,或通过提供两个续令牌来实现。
是否可以使用Java SDK实现该功能?或者有其他替代方案?希望能得到相关指导或建议。

我已完成如下实现,但根据Azure文档,LSN会标记到多个更新项。我知道可以使用setStartTime,但endTime该如何处理(如果可能,我需要获取更新到未来日期的文档)?或者能否通过对比两个续令牌来实现?

@Override
public Boolean call() throws Exception {
    ChangeFeedProcessorOptions options = new ChangeFeedProcessorOptions();
    AtomicReference<ChangeFeedProcessor> changeFeedProcessorRef = new AtomicReference<>();

    ChangeFeedProcessor changeFeedProcessor = new ChangeFeedProcessorBuilder()
            .hostName("Reprocess_LocalHost")
            .feedContainer(feedContainerId)
            .leaseContainer(leaseContainerId)
            .options(options)
            .handleChanges(changes -> {
                System.out.println("Start");
                for (JsonNode document : changes) {
                    long currentLsn = document.get("_lsn").asLong();
                    if (currentLsn <= toTokenLsn) {
                        System.out.println(" handleChanges() RECEIVED (within range): " + document.toPrettyString());
                    }
                    if (currentLsn >= toTokenLsn) {
                        System.out.println("Reached the specified toToken LSN: " + toTokenLsn + ". Stopping processing.");
                        ChangeFeedProcessor processorInstance = changeFeedProcessorRef.get();
                        if (processorInstance != null) {
                            processorInstance.stop()
                                    .subscribeOn(Schedulers.boundedElastic())
                                    .doOnSuccess(avoid -> System.out.println("Change Feed Processor has been stopped successfully."))
                                    .doOnError(error -> System.err.println("Failed to stop the Change Feed Processor: " + error))
                                    .subscribe();
                        }
                        return;
                    }
                }
                System.out.println("End");
            })
            .buildChangeFeedProcessor();

    changeFeedProcessorRef.set(changeFeedProcessor);
    changeFeedProcessor.start().block();
    return true;
}

方案指导

1. 日期范围筛选的实现

  • 起始时间:直接使用ChangeFeedProcessorOptions.setStartTime(Instant)设置变更的起始时间点,SDK会自动从该时间之后的变更开始处理。
  • 结束时间:SDK没有提供直接设置endTime的选项,可通过两种方式实现:
    • 在handleChanges回调中,检查每条变更文档的_ts字段(该字段是文档最后更新的Unix时间戳,单位秒),当文档时间戳超过目标结束时间时,标记停止处理器。
    • 如果需要持续处理未来新增的变更,无需设置结束条件,让处理器保持运行即可,它会自动监听并处理新的变更。

2. 基于续令牌的范围处理

续令牌代表变更源中的一个特定分页位置,要实现两个续令牌之间的变更处理:

  • 设置起始位置:使用ChangeFeedProcessorOptions.setStartContinuation(String)传入起始续令牌,处理器会从此位置开始处理。
  • 判断终止条件:处理过程中,通过FeedResponse.getContinuationToken()获取当前页的续令牌,当该令牌与目标结束续令牌匹配时,停止处理器。注意:续令牌是分页级别的,需确保完整处理当前页后再停止,避免漏处理文档。

3. 现有代码优化建议

当前通过LSN判断终止条件存在风险:一个LSN可能对应多个文档更新,容易出现提前停止或漏处理的情况。建议调整为:

  • 改用文档的_ts字段判断时间范围,更贴合业务的时间戳需求。
  • 若使用续令牌终止,需在回调中获取当前页的续令牌,对比后再停止,避免中断当前页的处理。

优化后代码示例(时间范围版)

@Override
public Boolean call() throws Exception {
    Instant startTime = Instant.parse("2024-01-01T00:00:00Z");
    Instant endTime = Instant.parse("2024-01-31T23:59:59Z");
    ChangeFeedProcessorOptions options = new ChangeFeedProcessorOptions()
            .setStartTime(startTime);
    AtomicReference<ChangeFeedProcessor> changeFeedProcessorRef = new AtomicReference<>();

    ChangeFeedProcessor changeFeedProcessor = new ChangeFeedProcessorBuilder()
            .hostName("Reprocess_LocalHost")
            .feedContainer(feedContainerId)
            .leaseContainer(leaseContainerId)
            .options(options)
            .handleChanges(changes -> {
                boolean shouldStop = false;
                System.out.println("Start");
                for (JsonNode document : changes) {
                    long documentTs = document.get("_ts").asLong();
                    Instant documentTime = Instant.ofEpochSecond(documentTs);
                    if (documentTime.isBefore(startTime)) {
                        continue;
                    }
                    if (documentTime.isAfter(endTime)) {
                        shouldStop = true;
                        break;
                    }
                    System.out.println(" handleChanges() RECEIVED (within range): " + document.toPrettyString());
                }
                if (shouldStop) {
                    System.out.println("Reached the specified end time: " + endTime + ". Stopping processing.");
                    ChangeFeedProcessor processorInstance = changeFeedProcessorRef.get();
                    if (processorInstance != null) {
                        processorInstance.stop()
                                .subscribeOn(Schedulers.boundedElastic())
                                .doOnSuccess(avoid -> System.out.println("Change Feed Processor has been stopped successfully."))
                                .doOnError(error -> System.err.println("Failed to stop the Change Feed Processor: " + error))
                                .subscribe();
                    }
                }
                System.out.println("End");
            })
            .buildChangeFeedProcessor();

    changeFeedProcessorRef.set(changeFeedProcessor);
    changeFeedProcessor.start().block();
    return true;
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 15:22:34