Flynn的服务发现组件(discoverd)部署在所有的Flynn集群节点上,该组件使用了Raft协议保证数据的一致性,目前采用的是hashicorp的实现(https://github.com/hashicorp/raft),我们知道在Raft协议中,如果参与选举的节点太多,会导致性能下降,那是不是说Flynn不支持大规模的节点呢?
Flynn是能够支持大规模节点的,虽然discoverd组件部署在所有的节点上,但并不是所有的节点都参与选举,只有部分节点作为Raft集群的节点,其余节点作为代理节点(proxying),Flynn通过以下逻辑判断是否是代理节点,启动discoverd时指定的-peers参数起了关键作用。
1 // if the advertise addr is not in the peer list we are proxying 2 proxying := true 3 for _, addr := range m.peers { 4 if addr == m.advertiseAddr { 5 proxying = false 6 break 7 } 8 }
discoverd组件本身依赖Raft协议保证数据的一致性,这里提到的数据,在discoverd里指的是服务,我们可以把我们的服务注册到discoverd上,我们的这些服务可能也需要一个Leader节点,例如Flynn的调度组件(scheduler)就只能在Leader节点上执行调度任务。那是否这些服务也是依赖Raft协议来选择Leader节点呢?
注册到discoverd组件上的服务不依赖Raft协议选择Leader,discoverd组件根据注册时间的长短来选择Leader,活的最长的服务节点被选择为Leader。所有注册到discoverd组件上的服务必须向discoverd组件发送心跳信息,discoverd组件如果检测到某个服务节点没有了心跳信息,就会把该节点移除,如果该节点恰好是Leader节点,那么就会触发重新选择Leader的动作。
以下代码是心跳检查相关代码
1// Open starts the raft consensus and opens the store. 2func (s *Store) Open() error { 3 go s.expirer() 4 5 return nil 6} 7 8// expirer runs in a separate goroutine and checks for instance expiration. 9func (s *Store) expirer() { 10 defer s.wg.Done() 11 12 ticker := time.NewTicker(s.ExpiryCheckInterval) 13 defer ticker.Stop() 14 15 for { 16 // Wait for next check or for close signal. 17 select { 18 case <-s.closing: 19 return 20 case <-ticker.C: 21 } 22 23 // Check all instances for expiration. 24 if err := s.EnforceExpiry(); err != nil && err != raft.ErrNotLeader { 25 s.logger.Printf("enforce expiry: %s", err) 26 } 27 } 28} 29 30// 以下代码有删减,仅保留主要逻辑 31// EnforceExpiry checks all instances for expiration and issues an expiration command, if necessary. 32// This function returns raft.ErrNotLeader if this store is not the current leader. 33func (s *Store) EnforceExpiry() error { 34 var cmd []byte 35 // Ignore if this store is not the leader and hasn't been for at least 2 TTLs intervals. 36 if !s.IsLeader() { 37 return raft.ErrNotLeader 38 } else if s.leaderTime.IsZero() || time.Since(s.leaderTime) < (2*s.InstanceTTL) { 39 return ErrLeaderWait 40 } 41 42 // Iterate over services and then instances. 43 var instances []expireInstance 44 for service, m := range s.data.Instances { 45 for _, inst := range m { 46 // Ignore instances that have heartbeated within the TTL. 47 if t := s.heartbeats[instanceKey{service, inst.ID}]; time.Since(t) <= s.InstanceTTL { 48 continue 49 } 50 51 // Add to list of instances to expire. 52 // The current expiry time is added to prevent a race condition of 53 // instances updating their expiry date while this command is applying. 54 instances = append(instances, expireInstance{ 55 Service: service, 56 InstanceID: inst.ID, 57 }) 58 } 59 } 60 61 // Create command to expire instances. 62 cmd, err := json.Marshal(&expireInstancesCommand{ 63 Instances: instances, 64 }) 65 66 // Apply command to raft. 67 if _, err := s.raftApply(expireInstancesCommandType, cmd); err != nil { 68 return err 69 } 70 return nil 71} 72
以下是出发重新选举的代码:
1func (s *Store) Apply(l *raft.Log) interface{} { 2 // Extract the command type and data. 3 typ, cmd := l.Data[0], l.Data[1:] 4 5 // Determine the command type by the first byte. 6 switch typ { 7 8 case expireInstancesCommandType: 9 return s.applyExpireInstancesCommand(cmd) 10 default: 11 return fmt.Errorf("invalid command type: %d", typ) 12 } 13} 14 15func (s *Store) applyExpireInstancesCommand(cmd []byte) error { 16 var c expireInstancesCommand 17 if err := json.Unmarshal(cmd, &c); err != nil { 18 return err 19 } 20 21 // Iterate over instances and remove ones with matching expiry times. 22 services := make(map[string]struct{}) 23 for _, expireInstance := range c.Instances { 24 // Remove instance. 25 delete(m, expireInstance.InstanceID) 26 27 // Broadcast down event. 28 s.broadcast(&discoverd.Event{ 29 Service: expireInstance.Service, 30 Kind: discoverd.EventKindDown, 31 Instance: inst, 32 }) 33 34 // Keep track of services invalidated. 35 services[expireInstance.Service] = struct{}{} 36 } 37 38 // Invalidate all services that had expirations. 39 for service := range services { 40 s.invalidateServiceLeader(service) 41 } 42 43 return nil 44} 45 46// invalidateServiceLeader updates the current leader of service. 47func (s *Store) invalidateServiceLeader(service string) { 48 // Retrieve service config. 49 c := s.data.Services[service] 50 51 // Ignore if there is no config or the leader is manually elected. 52 if c == nil || c.LeaderType == discoverd.LeaderTypeManual { 53 return 54 } 55 56 // Retrieve current leader ID. 57 prevLeaderID := s.data.Leaders[service] 58 59 // Find the oldest, non-expired instance. 60 var leader *discoverd.Instance 61 for _, inst := range s.data.Instances[service] { 62 if leader == nil || inst.Index < leader.Index { 63 leader = inst 64 } 65 } 66 67 // Retrieve the leader ID. 68 var leaderID string 69 if leader != nil { 70 leaderID = leader.ID 71 } 72 73 // Set leader. 74 s.data.Leaders[service] = leaderID 75 76 // Broadcast event. 77 if prevLeaderID != leaderID { 78 var inst *discoverd.Instance 79 if s.data.Instances[service] != nil { 80 inst = s.data.Instances[service][leaderID] 81 } 82 83 s.broadcast(&discoverd.Event{ 84 Service: service, 85 Kind: discoverd.EventKindLeader, 86 Instance: inst, 87 }) 88 } 89}
心跳信息是在哪里呢,估计都想到了,一定是在注册服务的那里。discoverd提供了对应的客户端,代码在discoverd/client里面,有个heartbeat.go就是专门来发心跳的。我们可以通过discoverd客户端的AddServiceAndRegister方法来完成服务注册功能
1func (c *Client) AddServiceAndRegister(service, addr string) (Heartbeater, error) { 2 if err := c.maybeAddService(service); err != nil { 3 return nil, err 4 } 5 return c.Register(service, addr) 6} 7 8func (c *Client) Register(service, addr string) (Heartbeater, error) { 9 return c.RegisterInstance(service, &Instance{Addr: addr}) 10} 11 12func (c *Client) RegisterInstance(service string, inst *Instance) (Heartbeater, error) { 13 h := newHeartbeater(c, service, inst) 14 err := runAttempts.Run(func() error { 15 firstErr := make(chan error) 16 go h.run(firstErr) 17 return <-firstErr 18 }) 19 if err != nil { 20 return nil, err 21 } 22 23 return h, nil 24} 25 26func (h *heartbeater) run(firstErr chan<- error) { 27 path := fmt.Sprintf("/services/%s/instances/%s", h.service, h.inst.ID) 28 register := func() error { 29 h.Lock() 30 defer h.Unlock() 31 return h.client().Put(path, h.inst, nil) 32 } 33 34 timer := time.NewTimer(nextHeartbeat()) 35 for { 36 select { 37 case <-timer.C: 38 if err := register(); err != nil { 39 timer.Reset(nextHeartbeatFailing()) 40 break 41 } 42 timer.Reset(nextHeartbeat()) 43 case <-h.stop: 44 h.client().Delete(path) 45 close(h.done) 46 return 47 } 48 } 49}
discoverd也是以Http协议提供服务,比如通过GET /services/abc/leader来获取abc服务的Leader节点,当然,也可以使用SSE协议来监听abc服务的Leader变化事件。Flynn的调度组件(scheduler)就是采用SSE协议监听Leader节点的变化。
1 r.PUT("/services/:service/leader", h.servePutLeader) 2 r.GET("/services/:service/leader", h.serveGetLeader)