Corda主流调用子流更新同一状态触发UnexpectedFlowEndException求助
问题:Corda主流调用子流更新同一LinearState时出现UnexpectedFlowEndException错误
我在Corda中编写了用于更新ContainerState的主流UpdateContainerStateFlow,尝试在该主流中调用另一个自定义子流UpdateContainerStateFlow2来更新同一ID的ContainerState属性(注:此为测试示例,实际可在主流内完成)。在构建主流交易前调用该子流后,出现如下错误:
错误日志
[WARN] 13:07:03,441 [Mock network] interceptors.DumpHistoryOnErrorInterceptor. - Flow [fb6a42fd-b9bf-45bb-ab98-2f2510235f91] error {fiber-id=10000008, flow-id=fb6a42fd-b9bf-45bb-ab98-2f2510235f91, invocation_id=38c29575-637a-459b-b8e9-81a960e93c7c, invocation_timestamp=2022-12-07T13:07:03.064Z, origin=O=Mock Company 1, L=London, C=GB, session_id=38c29575-637a-459b-b8e9-81a960e93c7c, session_timestamp=2022-12-07T13:07:03.064Z, thread-id=406, tx_id=EB1EA4EEA934C9C874461A6DDE35892057EBC472F405AF95241784CEDE3D6C92}
错误信息
net.corda.core.flows.UnexpectedFlowEndException: Counter-flow errored
at Received unexpected counter-flow exception from peer O=Mock Company 1, L=London, C=GB.() ~[?:?]
我疑惑子流是否应挂起主流并完成FinalityFlow的所有步骤,现附上相关代码,请求纠正理解误区并解决问题。
主流代码
@InitiatingFlow @StartableByRPC class UpdateContainerStateFlow(private val containerId: UUID, val status: String, val comments: String) : FlowLogic<SignedTransaction>() { override val progressTracker = ProgressTracker() @Suspendable override fun call(): SignedTransaction { val queryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId))) val currentContainer = serviceHub.vaultService.queryBy<ContainerState>(queryCriteria).states.single() val notary = serviceHub.networkMapCache.getNotary( CordaX500Name.parse("O=Notary,L=London,C=GB")) val updateContainer = currentContainer.state.data.copy( status = status, commentsFromStakeHolders = comments ) val abc = subFlow(UpdateContainerStateFlow2( containerId = containerId, status = "subflowOne", comments = "test" )) val builder = TransactionBuilder(notary) .addCommand(ContainerContract.Commands.Update(), listOf(updateContainer.logisticCompany.owningKey, updateContainer.partnerCompany.owningKey)) .addInputState(currentContainer) .addOutputState(updateContainer) // Step 4. Verify and sign it with our KeyPair. builder.verify(serviceHub) val ptx = serviceHub.signInitialTransaction(builder) val sessions = (updateContainer.participants - ourIdentity).map { initiateFlow(it) }.toSet() val stx = subFlow(CollectSignaturesFlow(ptx, sessions)) // Step 7. Assuming no exceptions, we can now finalise the transaction return subFlow(FinalityFlow(stx, sessions)) } } @InitiatedBy(UpdateContainerStateFlow::class) class UpdateContainerStateFlowResponder(val counterpartySession: FlowSession) : FlowLogic<SignedTransaction>() { @Suspendable override fun call(): SignedTransaction { val signTransactionFlow = object : SignTransactionFlow(counterpartySession) { override fun checkTransaction(stx: SignedTransaction) = requireThat { //Addition checks } } val txId = subFlow(signTransactionFlow).id return subFlow(ReceiveFinalityFlow(counterpartySession, expectedTxId = txId)) } }
UpdateContainerStateFlow2()的实现与上述代码完全一致。
核心问题分析
- 子流角色定位错误:将
UpdateContainerStateFlow2标记为@InitiatingFlow,作为子流调用时会尝试向参与者发起新会话,与主流已建立的会话冲突,触发异常。 - LinearState版本冲突:子流完成FinalityFlow后,目标State版本已升级,主流仍使用旧版本构建交易,会导致Vault状态校验失败。
- 子流设计误区:内部逻辑拆分的子流不应设计为完整的发起-响应式Flow,更不应单独提交FinalityFlow,否则会打断主流流程,造成节点状态不一致。
修正方案
方案一:将子流改为内部逻辑流(推荐)
如果只是拆分内部逻辑,将子流改为普通FlowLogic,去掉跨节点交互标记,仅负责状态预处理,不独立提交交易:
// 子流改为普通FlowLogic,仅处理状态更新逻辑 class UpdateContainerStateFlow2(private val containerId: UUID, val status: String, val comments: String) : FlowLogic<ContainerState>() { @Suspendable override fun call(): ContainerState { val queryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId))) val currentContainer = serviceHub.vaultService.queryBy<ContainerState>(queryCriteria).states.single() // 返回更新后的状态,由主流统一构建交易 return currentContainer.state.data.copy( status = status, commentsFromStakeHolders = comments ) } }
主流中调整调用逻辑:
// 调用子流获取处理后的状态 val subUpdatedContainer = subFlow(UpdateContainerStateFlow2( containerId = containerId, status = "subflowOne", comments = "test" )) // 使用子流返回的状态构建主流交易 val builder = TransactionBuilder(notary) .addCommand(ContainerContract.Commands.Update(), listOf(subUpdatedContainer.logisticCompany.owningKey, subUpdatedContainer.partnerCompany.owningKey)) .addInputState(currentContainer) .addOutputState(subUpdatedContainer)
方案二:子流作为独立交易(仅特殊场景使用)
如果必须让子流提交独立交易,主流需在调用后重新查询最新版本的State:
// 调用子流完成独立交易 val abc = subFlow(UpdateContainerStateFlow2( containerId = containerId, status = "subflowOne", comments = "test" )) // 重新查询最新版本的ContainerState val updatedQueryCriteria = QueryCriteria.LinearStateQueryCriteria(linearId = listOf(UniqueIdentifier(id = containerId))) val latestContainer = serviceHub.vaultService.queryBy<ContainerState>(updatedQueryCriteria).states.single() // 基于最新状态构建主流交易 val updateContainer = latestContainer.state.data.copy( status = status, commentsFromStakeHolders = comments ) val builder = TransactionBuilder(notary) .addCommand(ContainerContract.Commands.Update(), listOf(updateContainer.logisticCompany.owningKey, updateContainer.partnerCompany.owningKey)) .addInputState(latestContainer) .addOutputState(updateContainer)
关键原则总结
- 内部逻辑拆分的子流无需标记
@InitiatingFlow,不独立提交交易,由主流统一处理跨节点交互。 - 同一LinearState的连续更新尽量在单个交易中完成,避免跨交易版本冲突。
- 禁止嵌套调用
@InitiatingFlow,如需跨节点交互,由主流统一管理会话。
内容的提问来源于stack exchange,提问作者Benjamin
相关产品推荐
相关产品推荐

