【问题标题】:Corda RPC Connection Pooling/CachingCorda RPC 连接池/缓存
【发布时间】:2018-09-18 22:55:31
【问题描述】:

Corda 有连接池功能吗?如何处理多个 RPC 用户连接池...

如果您可以重定向到任何用于 RPC 连接池/缓存的开源实现/指南,请不胜感激...

【问题讨论】:

  • 你可以试试apache commons pool。我还没有用它来池化 RPC 连接,但我早就用它来池化 ssh 连接了。我相信它也应该适合你。

标签: rpc corda


【解决方案1】:

这是一个如何池化节点 RPC 连接的示例:

import net.corda.client.rpc.CordaRPCClient
import net.corda.client.rpc.CordaRPCConnection
import net.corda.core.utilities.NetworkHostAndPort
import net.corda.core.utilities.contextLogger
import net.corda.core.utilities.getOrThrow
import net.corda.node.services.Permissions
import net.corda.testing.driver.DriverParameters
import net.corda.testing.driver.NodeParameters
import net.corda.testing.driver.driver
import net.corda.testing.node.User
import org.junit.Test
import java.util.concurrent.ConcurrentHashMap
import java.util.concurrent.LinkedBlockingQueue

data class UserParams(val username: String, val password: String)

class PooledRpcConnections(address: NetworkHostAndPort): AutoCloseable {
    val userToPool = ConcurrentHashMap<UserParams, LinkedBlockingQueue<CordaRPCConnection>>()
    val client = CordaRPCClient(address)

    fun <A> withConnection(userParams: UserParams, block: (CordaRPCConnection) -> A): A {
        val queue = userToPool.getOrPut(userParams) { LinkedBlockingQueue() }
        val connection = queue.poll() ?: client.start(userParams.username, userParams.password)
        return try {
            block(connection)
        } finally {
            queue.add(connection)
        }
    }

    override fun close() {
        for (queue in userToPool.values) {
            do {
                val connection = queue.poll()
                connection?.close()
            } while (connection != null)
        }
    }
}

class Test {
    companion object {
        val log = contextLogger()
    }

    @Test
    fun poolWorks() {
        val users = ('a' .. 'f').map { User(it.toString(), it.toString(), setOf(Permissions.all())) }
        val userParams = users.map { UserParams(it.username, it.password) }
        driver(DriverParameters(startNodesInProcess = true, notarySpecs = emptyList())) {
            log.info("Starting node for users ${users.map { it.username }}")
            val node = startNode(NodeParameters(rpcUsers = users)).getOrThrow()
            log.info("Starting pool")
            PooledRpcConnections(node.rpcAddress).use { pool ->
                val N = 1000
                log.info("Making $N requests using pooled connections")
                (1 .. N).toList().parallelStream().forEach { i ->
                    val user = userParams[i % users.size]
                    pool.withConnection(user) { connection ->
                        log.info("USER[${user.username}] CONNECTION[${connection.hashCode()}] NODE_TIME[${connection.proxy.currentNodeTime()}]")
                    }
                }
                log.info("Done! Number of connections used per user: ${pool.userToPool.map { it.key.username to it.value.size }}")
            }
        }
    }
}

【讨论】:

    猜你喜欢
    • 2019-11-03
    • 2020-04-26
    • 2018-12-19
    • 1970-01-01
    • 2013-03-02
    • 1970-01-01
    • 1970-01-01
    • 2018-09-30
    • 2011-06-11
    相关资源
    最近更新 更多