能否在单个Corda流中执行多个流并创建多笔不同状态的交易?
在单个Corda流中执行多子流创建多交易的实现示例
当然可以!在Corda的流设计体系里,完全支持在一个主协调流中调用多个子流,分别生成针对不同状态的独立交易——这在需要一次性完成关联但不同的状态创建场景中非常实用。下面我给你一个基于Scala的完整示例,帮你快速理解并实现这个需求。
核心思路
主流负责整体流程的编排,每个子流专注于处理单一类型状态的全流程交易创建(包括构建交易结构、收集参与者签名、完成公证上链等步骤)。主流可以按顺序调用子流(适合有依赖关系的场景,比如付款需要关联发票ID),也可以通过async/await实现并行执行(适合无依赖的独立状态创建)。
代码示例
1. 定义两个示例状态
首先我们先定义两个不同的状态类型,比如发票状态和付款状态:
import net.corda.core.contracts.{ContractState, LinearState} import net.corda.core.identity.AbstractParty import net.corda.core.contracts.UniqueIdentifier // 发票状态 case class InvoiceState( amount: Int, issuer: AbstractParty, receiver: AbstractParty, override val participants: List[AbstractParty] ) extends LinearState { override val linearId: UniqueIdentifier = UniqueIdentifier() } // 付款状态 case class PaymentState( invoiceId: UniqueIdentifier, paidAmount: Int, payer: AbstractParty, payee: AbstractParty, override val participants: List[AbstractParty] ) extends LinearState { override val linearId: UniqueIdentifier = UniqueIdentifier() }
(注:实际使用时需要为每个状态绑定对应的Contract实现,这里简化了合约逻辑)
2. 实现子流:分别创建单个状态的交易
接下来编写两个子流,分别负责创建发票交易和付款交易:
import net.corda.core.flows.* import net.corda.core.transactions.TransactionBuilder import net.corda.core.contracts.Command // 创建发票的子流 @InitiatingFlow @StartableByRPC class CreateInvoiceFlow(private val amount: Int, private val receiver: AbstractParty) extends FlowLogic[SignedTransaction] { @Suspendable override def call(): SignedTransaction = { // 1. 获取自身身份和公证人 val me = ourIdentity val notary = serviceHub.networkMapCache.notaryIdentities.head // 2. 构建发票状态 val invoiceState = InvoiceState(amount, me, receiver, List(me, receiver)) // 3. 构建交易 val txBuilder = new TransactionBuilder(notary) .addOutputState(invoiceState, InvoiceContract.ID) .addCommand(new Command(InvoiceContract.Commands.Create(), List(me.getOwningKey))) // 4. 验证交易合法性 txBuilder.verify(serviceHub) // 5. 自身签名 val partSignedTx = serviceHub.signInitialTransaction(txBuilder) // 6. 收集对方签名并完成公证上链 val otherPartySession = initiateFlow(receiver) val fullySignedTx = subFlow(new CollectSignaturesFlow(partSignedTx, List(otherPartySession))) subFlow(new FinalityFlow(fullySignedTx, List(otherPartySession))) } } // 创建付款的子流 @InitiatingFlow @StartableByRPC class CreatePaymentFlow(private val invoiceId: UniqueIdentifier, private val paidAmount: Int, private val payee: AbstractParty) extends FlowLogic[SignedTransaction] { @Suspendable override def call(): SignedTransaction = { val me = ourIdentity val notary = serviceHub.networkMapCache.notaryIdentities.head val paymentState = PaymentState(invoiceId, paidAmount, me, payee, List(me, payee)) val txBuilder = new TransactionBuilder(notary) .addOutputState(paymentState, PaymentContract.ID) .addCommand(new Command(PaymentContract.Commands.Record(), List(me.getOwningKey))) txBuilder.verify(serviceHub) val partSignedTx = serviceHub.signInitialTransaction(txBuilder) val otherPartySession = initiateFlow(payee) val fullySignedTx = subFlow(new CollectSignaturesFlow(partSignedTx, List(otherPartySession))) subFlow(new FinalityFlow(fullySignedTx, List(otherPartySession))) } }
3. 实现主流:协调调用多个子流
最后编写主流,一次性调用上面两个子流,完成发票和付款状态的创建:
@InitiatingFlow @StartableByRPC class CreateInvoiceAndPaymentFlow( private val invoiceAmount: Int, private val receiver: AbstractParty, private val paidAmount: Int ) extends FlowLogic[(SignedTransaction, SignedTransaction)] { @Suspendable override def call(): (SignedTransaction, SignedTransaction) = { // 1. 先调用创建发票的子流,获取发票交易和状态ID val invoiceTx = subFlow(new CreateInvoiceFlow(invoiceAmount, receiver)) val invoiceState = invoiceTx.getTx.getOutputsOfType(classOf[InvoiceState]).head val invoiceId = invoiceState.linearId // 2. 再调用创建付款的子流,关联刚创建的发票ID val paymentTx = subFlow(new CreatePaymentFlow(invoiceId, paidAmount, receiver)) // 返回两个交易的结果 (invoiceTx, paymentTx) } }
关键注意事项
- 交易独立性:每个子流创建的是完全独立的交易,各自拥有独立的公证人、签名集合和状态数据,互相之间不会产生干扰
- 依赖处理:如果子流之间存在依赖(比如付款交易需要引用刚创建的发票ID),一定要按顺序调用子流,确保前序交易完成并获取到所需数据后再执行后续子流
- 并行优化:如果多个子流之间没有依赖关系,可以用Corda的
async函数实现并行执行,提升流程效率,示例代码如下:@Suspendable override def call(): (SignedTransaction, SignedTransaction) = { // 并行启动两个子流 val invoiceFuture = async { subFlow(new CreateInvoiceFlow(invoiceAmount, receiver)) } val paymentFuture = async { subFlow(new CreatePaymentFlow(preDefinedInvoiceId, paidAmount, receiver)) } // 等待两个子流完成并获取结果 val invoiceTx = invoiceFuture.get val paymentTx = paymentFuture.get (invoiceTx, paymentTx) } - 会话复用:如果多个子流需要和同一参与者交互,可以复用已建立的
FlowSession对象,避免重复建立网络连接,提升性能
内容的提问来源于stack exchange,提问作者scala
相关产品推荐
相关产品推荐

