mirror of
https://github.com/duanhf2012/origin.git
synced 2026-02-04 06:54:45 +08:00
优化代码
This commit is contained in:
@@ -9,9 +9,7 @@ import (
|
||||
"sync"
|
||||
)
|
||||
|
||||
var configdir = "./config/"
|
||||
|
||||
|
||||
var configDir = "./config/"
|
||||
|
||||
type NodeInfo struct {
|
||||
NodeId int
|
||||
@@ -21,83 +19,77 @@ type NodeInfo struct {
|
||||
}
|
||||
|
||||
type NodeRpcInfo struct {
|
||||
nodeinfo NodeInfo
|
||||
nodeInfo NodeInfo
|
||||
client *rpc.Client
|
||||
}
|
||||
|
||||
|
||||
var cluster Cluster
|
||||
|
||||
type Cluster struct {
|
||||
localNodeInfo NodeInfo //×
|
||||
localServiceCfg map[string]interface{} //map[servicename]配置数据*
|
||||
|
||||
mapRpc map[int] NodeRpcInfo//nodeid
|
||||
rpcServer rpc.Server
|
||||
serviceDiscovery IServiceDiscovery //服务发现接口
|
||||
|
||||
mapIdNode map[int]NodeInfo //map[NodeId]NodeInfo
|
||||
mapServiceNode map[string][]int //map[serviceName]NodeInfo
|
||||
|
||||
localNodeInfo NodeInfo
|
||||
localServiceCfg map[string]interface{} //map[serviceName]配置数据*
|
||||
mapRpc map[int] NodeRpcInfo //nodeId
|
||||
serviceDiscovery IServiceDiscovery //服务发现接口
|
||||
mapIdNode map[int]NodeInfo //map[NodeId]NodeInfo
|
||||
mapServiceNode map[string][]int //map[serviceName]NodeInfo
|
||||
locker sync.RWMutex
|
||||
rpcServer rpc.Server
|
||||
}
|
||||
|
||||
func SetConfigDir(cfgdir string){
|
||||
configdir = cfgdir
|
||||
func SetConfigDir(cfgDir string){
|
||||
configDir = cfgDir
|
||||
}
|
||||
|
||||
func SetServiceDiscovery(serviceDiscovery IServiceDiscovery) {
|
||||
cluster.serviceDiscovery = serviceDiscovery
|
||||
}
|
||||
|
||||
func (slf *Cluster) serviceDiscoveryDelNode (nodeId int){
|
||||
slf.locker.Lock()
|
||||
defer slf.locker.Unlock()
|
||||
func (cls *Cluster) serviceDiscoveryDelNode (nodeId int){
|
||||
cls.locker.Lock()
|
||||
defer cls.locker.Unlock()
|
||||
|
||||
slf.delNode(nodeId)
|
||||
cls.delNode(nodeId)
|
||||
}
|
||||
|
||||
func (slf *Cluster) delNode(nodeId int){
|
||||
func (cls *Cluster) delNode(nodeId int){
|
||||
//删除rpc连接关系
|
||||
rpc,ok := slf.mapRpc[nodeId]
|
||||
rpc,ok := cls.mapRpc[nodeId]
|
||||
if ok == true {
|
||||
delete(slf.mapRpc,nodeId)
|
||||
delete(cls.mapRpc,nodeId)
|
||||
rpc.client.Close(false)
|
||||
}
|
||||
|
||||
nodeInfo,ok := slf.mapIdNode[nodeId]
|
||||
nodeInfo,ok := cls.mapIdNode[nodeId]
|
||||
if ok == false {
|
||||
return
|
||||
}
|
||||
|
||||
for _,serviceName := range nodeInfo.ServiceList{
|
||||
slf.delServiceNode(serviceName,nodeId)
|
||||
cls.delServiceNode(serviceName,nodeId)
|
||||
}
|
||||
|
||||
delete(slf.mapIdNode,nodeId)
|
||||
delete(cls.mapIdNode,nodeId)
|
||||
}
|
||||
|
||||
func (slf *Cluster) delServiceNode(serviceName string,nodeId int){
|
||||
nodeList := slf.mapServiceNode[serviceName]
|
||||
func (cls *Cluster) delServiceNode(serviceName string,nodeId int){
|
||||
nodeList := cls.mapServiceNode[serviceName]
|
||||
for idx,nId := range nodeList {
|
||||
if nId == nodeId {
|
||||
slf.mapServiceNode[serviceName] = append(nodeList[idx:],nodeList[idx+1:]...)
|
||||
cls.mapServiceNode[serviceName] = append(nodeList[idx:],nodeList[idx+1:]...)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
func (slf *Cluster) serviceDiscoverySetNodeInfo (nodeInfo *NodeInfo){
|
||||
if nodeInfo.NodeId == slf.localNodeInfo.NodeId {
|
||||
func (cls *Cluster) serviceDiscoverySetNodeInfo (nodeInfo *NodeInfo){
|
||||
if nodeInfo.NodeId == cls.localNodeInfo.NodeId {
|
||||
return
|
||||
}
|
||||
|
||||
slf.locker.Lock()
|
||||
defer slf.locker.Unlock()
|
||||
cls.locker.Lock()
|
||||
defer cls.locker.Unlock()
|
||||
|
||||
//先清理删除
|
||||
slf.delNode(nodeInfo.NodeId)
|
||||
cls.delNode(nodeInfo.NodeId)
|
||||
|
||||
//再重新组装
|
||||
mapDuplicate := map[string]interface{}{} //预防重复数据
|
||||
@@ -108,58 +100,58 @@ func (slf *Cluster) serviceDiscoverySetNodeInfo (nodeInfo *NodeInfo){
|
||||
continue
|
||||
}
|
||||
|
||||
slf.mapServiceNode[serviceName] = append(slf.mapServiceNode[serviceName],nodeInfo.NodeId)
|
||||
cls.mapServiceNode[serviceName] = append(cls.mapServiceNode[serviceName],nodeInfo.NodeId)
|
||||
}
|
||||
|
||||
slf.mapIdNode[nodeInfo.NodeId] = *nodeInfo
|
||||
cls.mapIdNode[nodeInfo.NodeId] = *nodeInfo
|
||||
rpcInfo := NodeRpcInfo{}
|
||||
rpcInfo.nodeinfo = *nodeInfo
|
||||
rpcInfo.nodeInfo = *nodeInfo
|
||||
rpcInfo.client = &rpc.Client{}
|
||||
rpcInfo.client.Connect(nodeInfo.ListenAddr)
|
||||
slf.mapRpc[nodeInfo.NodeId] = rpcInfo
|
||||
cls.mapRpc[nodeInfo.NodeId] = rpcInfo
|
||||
}
|
||||
|
||||
func (slf *Cluster) buildLocalRpc(){
|
||||
func (cls *Cluster) buildLocalRpc(){
|
||||
rpcInfo := NodeRpcInfo{}
|
||||
rpcInfo.nodeinfo = slf.localNodeInfo
|
||||
rpcInfo.nodeInfo = cls.localNodeInfo
|
||||
rpcInfo.client = &rpc.Client{}
|
||||
rpcInfo.client.Connect("")
|
||||
|
||||
slf.mapRpc[slf.localNodeInfo.NodeId] = rpcInfo
|
||||
cls.mapRpc[cls.localNodeInfo.NodeId] = rpcInfo
|
||||
}
|
||||
|
||||
func (slf *Cluster) Init(localNodeId int) error{
|
||||
slf.locker.Lock()
|
||||
|
||||
func (cls *Cluster) Init(localNodeId int) error{
|
||||
cls.locker.Lock()
|
||||
|
||||
//1.处理服务发现接口
|
||||
if slf.serviceDiscovery == nil {
|
||||
slf.serviceDiscovery = &ConfigDiscovery{}
|
||||
if cls.serviceDiscovery == nil {
|
||||
cls.serviceDiscovery = &ConfigDiscovery{}
|
||||
}
|
||||
|
||||
//2.初始化配置
|
||||
err := slf.InitCfg(localNodeId)
|
||||
err := cls.InitCfg(localNodeId)
|
||||
if err != nil {
|
||||
slf.locker.Unlock()
|
||||
cls.locker.Unlock()
|
||||
return err
|
||||
}
|
||||
|
||||
slf.rpcServer.Init(slf)
|
||||
slf.buildLocalRpc()
|
||||
cls.rpcServer.Init(cls)
|
||||
cls.buildLocalRpc()
|
||||
|
||||
slf.serviceDiscovery.RegFunDelNode(slf.serviceDiscoveryDelNode)
|
||||
slf.serviceDiscovery.RegFunSetNode(slf.serviceDiscoverySetNodeInfo)
|
||||
slf.locker.Unlock()
|
||||
cls.serviceDiscovery.RegFunDelNode(cls.serviceDiscoveryDelNode)
|
||||
cls.serviceDiscovery.RegFunSetNode(cls.serviceDiscoverySetNodeInfo)
|
||||
cls.locker.Unlock()
|
||||
|
||||
err = slf.serviceDiscovery.Init(localNodeId)
|
||||
err = cls.serviceDiscovery.Init(localNodeId)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (slf *Cluster) FindRpcHandler(servicename string) rpc.IRpcHandler {
|
||||
pService := service.GetService(servicename)
|
||||
func (cls *Cluster) FindRpcHandler(serviceName string) rpc.IRpcHandler {
|
||||
pService := service.GetService(serviceName)
|
||||
if pService == nil {
|
||||
return nil
|
||||
}
|
||||
@@ -167,22 +159,22 @@ func (slf *Cluster) FindRpcHandler(servicename string) rpc.IRpcHandler {
|
||||
return pService.GetRpcHandler()
|
||||
}
|
||||
|
||||
func (slf *Cluster) Start() {
|
||||
slf.rpcServer.Start(slf.localNodeInfo.ListenAddr)
|
||||
func (cls *Cluster) Start() {
|
||||
cls.rpcServer.Start(cls.localNodeInfo.ListenAddr)
|
||||
}
|
||||
|
||||
func (slf *Cluster) Stop() {
|
||||
slf.serviceDiscovery.OnNodeStop()
|
||||
func (cls *Cluster) Stop() {
|
||||
cls.serviceDiscovery.OnNodeStop()
|
||||
}
|
||||
|
||||
func GetCluster() *Cluster{
|
||||
return &cluster
|
||||
}
|
||||
|
||||
func (slf *Cluster) GetRpcClient(nodeid int) *rpc.Client {
|
||||
slf.locker.RLock()
|
||||
defer slf.locker.RUnlock()
|
||||
c,ok := slf.mapRpc[nodeid]
|
||||
func (cls *Cluster) GetRpcClient(nodeId int) *rpc.Client {
|
||||
cls.locker.RLock()
|
||||
defer cls.locker.RUnlock()
|
||||
c,ok := cls.mapRpc[nodeId]
|
||||
if ok == false {
|
||||
return nil
|
||||
}
|
||||
@@ -205,7 +197,7 @@ func GetRpcClient(nodeId int,serviceMethod string,clientList *[]*rpc.Client) err
|
||||
return fmt.Errorf("servicemethod param %s is error!",serviceMethod)
|
||||
}
|
||||
|
||||
//1.找到对应的rpcnodeid
|
||||
//1.找到对应的rpcNodeid
|
||||
GetCluster().GetNodeIdByService(serviceAndMethod[0],clientList)
|
||||
return nil
|
||||
}
|
||||
@@ -214,7 +206,7 @@ func GetRpcServer() *rpc.Server{
|
||||
return &cluster.rpcServer
|
||||
}
|
||||
|
||||
func (slf *Cluster) IsNodeConnected (nodeId int) bool {
|
||||
pClient := slf.GetRpcClient(nodeId)
|
||||
func (cls *Cluster) IsNodeConnected (nodeId int) bool {
|
||||
pClient := cls.GetRpcClient(nodeId)
|
||||
return pClient!=nil && pClient.IsConnected()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user