【发布时间】:2021-05-21 05:59:31
【问题描述】:
我正在使用 Java 8,并且我正在尝试运行一系列 CompletionStage。
我不想使用join() 或get(),我想明确地完成CompletionStage。
我正在尝试运行两个数据库查询,第二个依赖于第一个查询的结果。我正在使用会话启动数据库事务,运行 write query1,write query2 并且只有当两者都成功时,我才想提交事务或将其回滚。 事务和会话是 Neo4j java API https://neo4j.com/docs/api/java-driver/current/org/neo4j/driver/async/AsyncSession.html#writeTransactionAsync-org.neo4j.driver.async.AsyncTransactionWork-
的一部分在运行两个查询成功/失败后,我想关闭会话(标准数据库实践)
这是伪代码-
DB Session starts transaction
run Write Query1
run Write Query2
if both are successful
commit transaction
else
rollback transaction
close session
我想要实现的是,如果 query1/query2 失败,那么它应该只是回滚事务并关闭会话。
如果 Query1 的结果不正确(小于某个阈值),查询 1 也可以抛出 CustomException。在这种情况下,它应该回滚事务。我正在为每个查询回滚 exceptionally 块中的事务。
快乐路径在下面的代码中运行良好,但是当我想抛出 CustomException 时,不会调用 Query2 块,甚至永远不会调用 Completable.allOf。
CompletableFuture<String> firstFuture = new CompletableFuture();
CompletableFuture<String> secondFuture = new CompletableFuture();
CompletableFuture<String> lastFuture = new CompletableFuture();
//Lambda that executes transaction
TransactionWork<CompletionStage<String>> runTransactionWork = transaction -> {
//Write Query1
transaction.runAsync("DB WRITE QUERY1") //Running Write Query 1
.thenCompose(someFunctionThatReturnsCompletionStage)
.thenApply(val -> {
//throw CustomException if value less then threshold
if(val < threshold){
throw new CustomException("Incorrect value found");
}else{
//if value is correct then complete future
firstFuture.complete(val);
}
firstQuery.complete(val);
}).exceptionally(error -> {
//Since failure occured in Query1 want to roll back
transaction.rollbackAsync();
firstFuture.completeExceptionally(error);
throw new RuntimeException("There has been an error in first query " + error.getMessage());
});
//after the first write query is done then run the second write query
firstFuture.thenCompose(val -> transaction.runAsync("DB Write QUERY 2"))
.thenCompose(someFunctionThatReturnsCompletionStage)
.thenApply(val -> {
//if value is correct then complete
secondFuture.complete(val);
}
}).exceptionally(error -> {
//Incase of failure in Query2 want to roll back
transaction.rollbackAsync();
secondFuture.completeExceptionally(error);
throw new RuntimeException("There has been an error in second query " + error.getMessage());
});
//wait for both to complete and then complete the last future
CompletableFuture.allOf(firstFuture, secondFuture)
.handle((empty, ex) -> {
if(ex != null){
lastFuture.completeExceptionally(ex);
}else{
//commit the transaction
transaction.commitAsync();
lastFuture.complete("OK");
}
return lastFuture;
});
return lastFuture;
}
//Create a database session
Session session = driver.session();
//runTransactionWork is lambda that has access to transaction
session.writeTransactionAsync(runTransactionWork)
.handle((val, err) -> {
if(val != null){
session.closeAsync();
//send message to some broker about success
}else{
//fail logic
}
});
如何实现短路异常以确保事务回滚并直接进入会话异常块。
这些是我对基于不同用例调用的代码块的观察,注意这些是基于我在代码中放置的调试点 -
- 快乐的路径 - firstFuture(success) -> secondFuture(success) -> LastFuture (success) -> 会话块成功调用(工作正常)
- 第一个 Future 失败 - firstFuture(由于异常而失败)-> secondFuture(从未调用)-> LastFuture(从未调用)-> 会话块失败(从未调用)
- 第二个未来失败 - firstFuture(成功)-> secondFuture(由于异常而失败)-> LastFuture(从未调用)-> 会话块失败(从未调用)
我希望 #2 和 #3 也能正常工作,并且应该回滚相应的事务并关闭会话。
我的问题是,为什么当未来 completesExceptionally 之一时,allOf 的 handle 的例外部分不会被调用?
【问题讨论】:
-
你能在这里提供一个简化的例子吗?我已经读了几次这个问题(加上只对你有意义的代码),我不明白你的问题。
-
@Eugene 我做了一些编辑,还添加了一些伪代码。如果现在可以理解,请告诉我。感谢您的宝贵时间。
-
这里有点矛盾。一方面你说你只想在两者都成功时提交,但另一方面,你的观点 (2) 说你仍然想要执行第二个事务。这是为什么?如果第一个失败了,为什么要执行第二个?你还是回滚。你能澄清一下吗?
-
@Eugene 是的,我仍然想回滚,我正在本地调试以检查它是否达到调试点。所以这就是我添加(从未调用)的原因。你是对的,我不想调用第二个块。它应该只是回滚。
标签: java asynchronous exception java-8 completable-future