Skip to content

Commit 985d4de

Browse files
committed
Use replica status for throttling via control replicas in move-tables mode
1 parent 3546b11 commit 985d4de

2 files changed

Lines changed: 34 additions & 0 deletions

File tree

go/logic/throttler.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -209,6 +209,14 @@ func (thlr *Throttler) collectControlReplicasLag() {
209209
return lag, err
210210
}
211211

212+
if thlr.migrationContext.IsMoveTablesMode() {
213+
dbVersion, err := mysql.GetDBVersion(thlr.migrationContext.Uuid, dbUri)
214+
if err != nil {
215+
return lag, err
216+
}
217+
return mysql.GetReplicationLagFromSlaveStatus(dbVersion, db)
218+
}
219+
212220
if err := db.QueryRow(replicationLagQuery).Scan(&heartbeatValue); err != nil {
213221
return lag, err
214222
}

go/mysql/utils.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,8 @@ func (rlg *ReplicationLagResult) HasLag() bool {
4747
// knownDBs is a DB cache by uri
4848
var knownDBs map[string]*gosql.DB = make(map[string]*gosql.DB)
4949
var knownDBsMutex = &sync.Mutex{}
50+
var knownDBsVersions map[string]string = make(map[string]string)
51+
var knownDBsVersionsMutex = &sync.Mutex{}
5052

5153
func GetDB(migrationUuid string, mysql_uri string) (db *gosql.DB, exists bool, err error) {
5254
cacheKey := migrationUuid + ":" + mysql_uri
@@ -66,6 +68,30 @@ func GetDB(migrationUuid string, mysql_uri string) (db *gosql.DB, exists bool, e
6668
return db, exists, nil
6769
}
6870

71+
// GetDBVersion returns the MySQL version for a given mysql_uri, and caches it for future calls.
72+
// Uses GetDB to get a connection to the database.
73+
func GetDBVersion(migrationUuid string, mysql_uri string) (dbVersion string, err error) {
74+
cacheKey := migrationUuid + ":" + mysql_uri
75+
76+
knownDBsVersionsMutex.Lock()
77+
defer knownDBsVersionsMutex.Unlock()
78+
79+
if dbVersion, exists := knownDBsVersions[cacheKey]; exists {
80+
return dbVersion, nil
81+
}
82+
83+
db, _, err := GetDB(migrationUuid, mysql_uri)
84+
if err != nil {
85+
return "", err
86+
}
87+
var version string
88+
if err := db.QueryRow(`select @@global.version`).Scan(&version); err != nil {
89+
return "", err
90+
}
91+
knownDBsVersions[cacheKey] = version
92+
return version, nil
93+
}
94+
6995
// GetReplicationLagFromSlaveStatus returns replication lag for a given db; via SHOW SLAVE STATUS
7096
func GetReplicationLagFromSlaveStatus(dbVersion string, informationSchemaDb *gosql.DB) (replicationLag time.Duration, err error) {
7197
showReplicaStatusQuery := fmt.Sprintf("show %s", ReplicaTermFor(dbVersion, `slave status`))

0 commit comments

Comments
 (0)