Couchbase TestContainer事务场景下数据写入后查询不到的问题求助
Couchbase TestContainer事务场景下数据写入后查询不到的问题求助
大家好,我现在遇到了一个Couchbase TestContainer的异常问题:在集成测试中,我的业务代码通过Couchbase事务成功插入了数据(从代码日志可以看到插入成功的记录),但在测试代码里去查询这些数据时,却完全查不到结果,日志里显示查询返回的都是空集合。
核心业务代码说明
我的业务逻辑是通过Couchbase事务来处理数据插入,所有数据库操作都依赖TransactionAttemptContext来保证事务性:
事务执行主逻辑
public void execute() { boolean success = false; int retryCount = 0; while (!success) { try { transactions.run((TransactionAttemptContext ctx) -> { process(ctx); }); success = true; } catch (TransactionFailedException ex) { Throwable cause = ex.getCause(); if (cause instanceof DocumentExistsException || cause instanceof DuplicateKeyException) { if (cause.getMessage().contains(getCollectionName(TransactionSubLedgerEntity.class))) { log.warn("Duplicate event with eventId {} while stock-processing. Event already present", storeStockMovement.getStockTransaction().stockTransactionId()); success = true; } else { retryCount++; log.warn("Duplicate entry in stockWacSubLedger or stockRefSubLedger eventId {} while stock-processing. Retrying... Attempt {}", storeStockMovement.getStockTransaction().stockTransactionId(), retryCount); if (retryCount >= MAX_RETRY_COUNT) { throw new TransactionProcessingException( String.format("Exceeded max retries while processing srs event with id %s", storeStockMovement.getStockTransaction().stockTransactionId()), ex); } } } else if (cause instanceof WacCalculationTransientException) { throw (WacCalculationTransientException) cause; } else { throw new TransactionProcessingException( String.format("Error while processing stock event with id %s", storeStockMovement.getStockTransaction().stockTransactionId()), ex); } } } }
事务内的插入操作
所有数据写入都是通过事务上下文ctx完成的,比如这个保存实体的方法:
public TransactionSubLedgerEntity saveEntity(TransactionAttemptContext ctx, TransactionSubLedgerEntity transactionSubLedgerEntity) { Collection collection = getCollection(); try { ctx.insert(collection, transactionSubLedgerEntity.getTransactionId(), transactionSubLedgerEntity); log.info("Inserted new TransactionSubLedgerEntity with ID: {}", transactionSubLedgerEntity.getTransactionId()); } catch (DocumentExistsException e) { log.error("Duplicate TransactionSubLedgerEntity detected with ID: {}", transactionSubLedgerEntity.getTransactionId()); throw e; } return transactionSubLedgerEntity; }
测试代码及问题现象
我的测试代码是用Groovy写的,执行完业务逻辑后,我尝试用couchbaseTemplate和Repository两种方式查询数据,但结果都是空的,哪怕我已经sleep了1秒等待事务提交:
def "Verify that successful stock receipt processing txn"() { given: "A SRS transfer is present in the system" buildDepotStockMovement(tpnb, storeId, referenceId, srsTransactionDate) insertStockWac(tpnb, storeId, 10) if(rcvEventTpnb!=tpnb) insertStockWac(rcvEventTpnb,storeId,11) insertStockRef(referenceId, 10, tpnb, storeId, stockInTransit, createdAt, srsTransactionDate) //currency mock ConfigData configData = new ConfigData() configData.setCurrency("GBP") CurrencyMapping currencyMapping = new CurrencyMapping() def storeStockMovement = createStockReceivingEvent(rcvEventTpnb == null ? tpnb : rcvEventTpnb, storeId, invoiceNo, receiveQty, sourceTransactionDateTime) currencyMapping.setConfigData(Collections.singletonList(configData)) Mockito.when(configurationService.getCurrencyMappingConfig(STOCK_CURRENCY_COUNTRY_MAPPING, storeStockMovement.getLocationReferenceEnriched().country())) .thenReturn(currencyMapping) Mockito.when(wacService.fetchCurrentWac(Mockito.anyString(), Mockito.any(LocalDate.class), Mockito.anyString())) .thenReturn(BigDecimal.valueOf(15.25)) Mockito.when(missingPreReqHelper.persistMissingPreReqEvent()).thenReturn(new MissingPreReqEntity()) when: "Stock receipt is received for that referenceId" String auditId = "test-1" def command = new StockReceivingProcessingCommand(storeStockMovement, stockProcessorHelper, srsProcessorHelper, auditId, missingPreReqHelper, transactions) command.execute() Thread.sleep(1000) then: "Expected transaction TransactionCode=#transactionCode ReasonCode=#reasonCode is created" if (transactionCode >= 0 && reasonCode >= 0) { QueryCriteria criteria = QueryCriteria.where("storeId").is(storeId) .and(QueryCriteria.where("tpnb").is(rcvEventTpnb == null ? tpnb : rcvEventTpnb)) .and(QueryCriteria.where("referenceId").is(invoiceNo)) .and(QueryCriteria.where("transactionCode").is(transactionCode)) .and(QueryCriteria.where("reasonCode").is(reasonCode)) def results = couchbaseTemplate.findByQuery(TransactionSubLedgerEntity) .withConsistency(QueryScanConsistency.REQUEST_PLUS) .matching(criteria) .all() log.info("results {}", results); log.info("transaction {}", transactionSubLedgerRepository.findAll().toString()) assert results.size() == 1 assert results.stream().findFirst().get().quantity == expectedQty assert results.stream().findFirst().get().newStockOnHand == expectedSOH } }
执行测试时,日志里能看到业务代码输出的Inserted new TransactionSubLedgerEntity with ID: xxx,但测试代码里的results和transactionSubLedgerRepository.findAll()返回的都是空集合,导致断言失败。
有没有大佬遇到过类似的问题?麻烦帮忙分析下可能的原因,谢谢!
备注:内容来源于stack exchange,提问作者arqam
相关产品推荐
相关产品推荐

