func(bo*BackendObserver)checkHealth(ctxcontext.Context,backendsmap[string]*BackendInfo)map[string]*backendHealth{curBackendHealth:=make(map[string]*backendHealth,len(backends))foraddr,info:=rangebackends{bh:=&backendHealth{status:StatusHealthy,}curBackendHealth[addr]=bh// http 服务检查ifinfo!=nil&&len(info.IP)>0{schema:="http"httpCli:=*bo.httpClihttpCli.Timeout=bo.healthCheckConfig.DialTimeouturl:=fmt.Sprintf("%s://%s:%d%s",schema,info.IP,info.StatusPort,statusPathSuffix)resp,err:=httpCli.Get(url)iferr!=nil{bh.status=StatusCannotConnectbh.pingErr=errors.Wrapf(err,"connect status port failed")continue}}// tcp 服务检查conn,err:=net.DialTimeout("tcp",addr,bo.healthCheckConfig.DialTimeout)iferr!=nil{bh.status=StatusCannotConnectbh.pingErr=errors.Wrapf(err,"connect sql port failed")}}returncurBackendHealth}
// notifyIfChanged 根据最新的 tidb 拓扑 bhMap 与之前的 tidb 拓扑 bo.curBackendInfo 进行比较// - 在 bo.curBackendInfo 中但是不在 bhMap 中:说明 tidb 节点失联,需要记录下// - 在 bo.curBackendInfo 中也在 bhMap 中,但是最新的状态不是 StatusHealthy:也需要记录下// - 在 bhMap 中但是不在 bo.curBackendInfo 中:说明是新增 tidb 节点,需要记录下func(bo*BackendObserver)notifyIfChanged(bhMapmap[string]*backendHealth){updatedBackends:=make(map[string]*backendHealth)foraddr,lastHealth:=rangebo.curBackendInfo{iflastHealth.status==StatusHealthy{ifnewHealth,ok:=bhMap[addr];!ok{updatedBackends[addr]=&backendHealth{status:StatusCannotConnect,pingErr:errors.New("removed from backend list"),}updateBackendStatusMetrics(addr,lastHealth.status,StatusCannotConnect)}elseifnewHealth.status!=StatusHealthy{updatedBackends[addr]=newHealthupdateBackendStatusMetrics(addr,lastHealth.status,newHealth.status)}}}foraddr,newHealth:=rangebhMap{ifnewHealth.status==StatusHealthy{lastHealth,ok:=bo.curBackendInfo[addr]if!ok{lastHealth=&backendHealth{status:StatusCannotConnect,}}iflastHealth.status!=StatusHealthy{updatedBackends[addr]=newHealthupdateBackendStatusMetrics(addr,lastHealth.status,newHealth.status)}elseiflastHealth.serverVersion!=newHealth.serverVersion{// Not possible here: the backend finishes upgrading between two health checks.updatedBackends[addr]=newHealth}}}// Notify it even when the updatedBackends is empty, in order to clear the last error.bo.eventReceiver.OnBackendChanged(updatedBackends,nil)bo.curBackendInfo=bhMap}
typeScoreBasedRouterstruct{sync.Mutex// A list of *backendWrapper. The backends are in descending order of scores.backends*glist.List[*backendWrapper]// ...}// 被 BackendObserver 调用,传来的 backends 会合并到 ScoreBasedRouter::backends 中func(router*ScoreBasedRouter)OnBackendChanged(backendsmap[string]*backendHealth,errerror){}// 通过比较 backend 分数方式调整 ScoreBasedRouter::backends 中的位置func(router*ScoreBasedRouter)adjustBackendList(be*glist.Element[*backendWrapper]){}// 协程方式运行,做负载均衡处理func(router*ScoreBasedRouter)rebalanceLoop(ctxcontext.Context){}
typeBackendConnManagerstruct{// processLock makes redirecting and command processing exclusive.processLocksync.MutexclientIO*pnet.PacketIObackendIOatomic.Pointer[pnet.PacketIO]authenticator*Authenticator}func(mgr*BackendConnManager)Redirect(newAddrstring)bool{}func(mgr*BackendConnManager)processSignals(ctxcontext.Context){}func(mgr*BackendConnManager)tryRedirect(ctxcontext.Context){}func(mgr*BackendConnManager)querySessionStates(backendIO*pnet.PacketIO)(sessionStates,sessionTokenstring,errerror){}func(mgr*BackendConnManager)ExecuteCmd(ctxcontext.Context,request[]byte)(errerror){}
迁移消息接收
在前文的 rebalance 方法最后,有行这样的逻辑
conn.Redirect(idlestBackend.addr)
这就是 ScoreBasedRouter 的通知给对应 conn 的地方。
这里调用的是 BackendConnManager::Redirect, 具体执行逻辑
将目标 backend 存储到 redirectInfo
给 signalReceived channel 发 signalTypeRedirect 消息
func(mgr*BackendConnManager)Redirect(newAddrstring)bool{// NOTE: BackendConnManager may be closing concurrently because of no lock.switchmgr.closeStatus.Load(){casestatusNotifyClose,statusClosing,statusClosed:returnfalse}mgr.redirectInfo.Store(&signalRedirect{newAddr:newAddr})// Generally, it won't wait because the caller won't send another signal before the previous one finishes.mgr.signalReceived<-signalTypeRedirectreturntrue}
该消息被 BackendConnManager::processSignals 协程接收
func(mgr*BackendConnManager)processSignals(ctxcontext.Context){for{select{cases:=<-mgr.signalReceived:// Redirect the session immediately just in case the session is finishedTxn.mgr.processLock.Lock()switchs{casesignalTypeGracefulClose:mgr.tryGracefulClose(ctx)casesignalTypeRedirect:// <<<<<<<<<<<<<<<<<<mgr.tryRedirect(ctx)}mgr.processLock.Unlock()casers:=<-mgr.redirectResCh:mgr.notifyRedirectResult(ctx,rs)case<-mgr.checkBackendTicker.C:mgr.checkBackendActive()case<-ctx.Done():return}}}
func(cp*CmdProcessor)finishedTxn()bool{ifcp.serverStatus&(StatusInTrans|StatusQuit)>0{returnfalse}// If any result of the prepared statements is not fetched, we should wait.return!cp.hasPendingPreparedStmts()}func(cp*CmdProcessor)updatePrepStmtStatus(request[]byte,serverStatusuint16){var(stmtIDintprepStmtStatusuint32)cmd:=pnet.Command(request[0])switchcmd{casepnet.ComStmtSendLongData,pnet.ComStmtExecute,pnet.ComStmtFetch,pnet.ComStmtReset,pnet.ComStmtClose:stmtID=int(binary.LittleEndian.Uint32(request[1:5]))casepnet.ComResetConnection,pnet.ComChangeUser:cp.preparedStmtStatus=make(map[int]uint32)returndefault:return}switchcmd{casepnet.ComStmtSendLongData:prepStmtStatus=StatusPrepareWaitExecutecasepnet.ComStmtExecute:ifserverStatus&mysql.ServerStatusCursorExists>0{prepStmtStatus=StatusPrepareWaitFetch}casepnet.ComStmtFetch:ifserverStatus&mysql.ServerStatusLastRowSend==0{prepStmtStatus=StatusPrepareWaitFetch}}ifprepStmtStatus>0{cp.preparedStmtStatus[stmtID]=prepStmtStatus}else{delete(cp.preparedStmtStatus,stmtID)}}