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
相关产品推荐
相关产品推荐

