Vitess全局唯一ID生成的实现方案
在Vitess实现方案中,每个设置了全局唯一列的表,都会对应一张sequence序列表。例如对于表user,会对应一张名为user_seq的序列表,原表与序列表的关联关系会记录在元数据中。user表以及user_seq这两张表元数据信息分别如下:
user表元数据:分片键为name列,分片算法为hash;全局唯一列为id列,依赖user_seq表生成具体的值。
{"tables": {"user": {"column_vindexes": [{"column": "name","name": "hash"}],"auto_increment": {"column": "id","sequence": "user_seq"}}}}
user_seq表元数据:表类型标识为sequence。
{"tables": {"user_seq": {"type": "sequence"}}}
CREATE TABLE user_seq (id int,next_id bigint,cache bigint,PRIMARY KEY (id)) COMMENT 'vitess_sequence';
mysql> select * from user_seq;+----+---------+-------+| id | next_id | cache |+----+---------+-------+| 0 | 1000 | 100 |+----+---------+-------+
// 获取sequence的方法func (qre *QueryExecutor) execNextval() (*sqltypes.Result, error) {// 从plan中获取inc(为要获取的id数量)以及tableNameinc, err := resolveNumber(qre.plan.NextCount, qre.bindVars)tableName := qre.plan.TableName()t := qre.plan.Tablet.SequenceInfo.Lock()defer t.SequenceInfo.Unlock()if t.SequenceInfo.NextVal == 0 || t.SequenceInfo.NextVal+inc > t.SequenceInfo.LastVal {// 在事务中运行_, err := qre.execAsTransaction(func(conn *StatefulConnection) (*sqltypes.Result, error) {// 使用select for update锁住行数据以免在计算并更新新值期间被其他线程修改query := fmt.Sprintf("select next_id, cache from %s where id = 0 for update", sqlparser.String(tableName))qr, err := qre.execSQL(conn, query, false)nextID, err := evalengine.ToInt64(qr.Rows[0][0])if t.SequenceInfo.LastVal != nextID {// 如果从_seq表读取得到的id值小于tablet缓存中id,则将缓存中的值更新到_seq表中if nextID < t.SequenceInfo.LastVal {log.Warningf("Sequence next ID value %v is below the currently cached max %v, updating it to max", nextID, t.SequenceInfo.LastVal)nextID = t.SequenceInfo.LastVal}t.SequenceInfo.NextVal = nextIDt.SequenceInfo.LastVal = nextID}cache, err := evalengine.ToInt64(qr.Rows[0][1])// 按照cache的倍数获取到大于inc量的缓存,计算出新newLastnewLast := nextID + cachefor newLast < t.SequenceInfo.NextVal+inc {newLast += cache}// 将新的边界值更新到_seq表中query = fmt.Sprintf("update %s set next_id = %d where id = 0", sqlparser.String(tableName), newLast)_, err = qre.execSQL(conn, query, false)t.SequenceInfo.LastVal = newLast})}// 返回获取的sequence值 更新SequenceInforet := t.SequenceInfo.NextValt.SequenceInfo.NextVal += increturn ret}
public long[] querySequenceValue(Vcursor vCursor, ResolvedShard resolvedShard, String sequenceTableName) throws SQLException, InterruptedException {// cas 重试次数限制int retryTimes = DEFAULT_RETRY_TIMES;while (retryTimes > 0) {// 查询_seq表中的sequence设置,其中cache为本地缓存的大小String querySql = "select next_id, cache from " + sequenceTableName + " where id = 0";VtResultSet vtResultSet = (VtResultSet) vCursor.executeStandalone(querySql, new HashMap<>(), resolvedShard, false);long[] sequenceInfo = getVtResultValue(vtResultSet);long next = sequenceInfo[0];long cache = sequenceInfo[1];// 将计算出的next_id的值尝试更新到_seq表中,如果失败则重新读取并更新,直到成功为止String updateSql = "update " + sequenceTableName + " set next_id = " + (next + cache) + " where next_id =" + sequenceInfo[0];VtRowList vtRowList = vCursor.executeStandalone(updateSql, new HashMap<>(), resolvedShard, false);if (vtRowList.getRowsAffected() == 1) {sequenceInfo[0] = next;return sequenceInfo;}retryTimes--;Thread.sleep(ThreadLocalRandom.current().nextInt(1, 6));}throw new SQLException("Update sequence cache failed within retryTimes: " + DEFAULT_RETRY_TIMES);}
在整个查询并更新序列表的过程中,没有出现Vitess实现中的开启事务以及产生锁表的情况,而是使用了CAS更新的方式。 利用update user_seq set next_id=?where next_id=?执行的返回值判断是否语句是否更新成功,如果失败则重新查询next_id的值,计算新值再尝试更新, 如果出现并发争抢的情况,Vtdriver中允许最多的重试次数DEFAULT_RETRY_TIMES为100次。
在Vitess实现sequence的源码当中,其更新序列表的过程为:开启事务时执行select for update,使用表锁,保证多线程安全。在现实往往充满了不确定性,我们可以想象一下:如果应用锁了数据库中的表后,由于自身的性能原因等而迟迟没有执行commit操作,或者应用节点出现了宕机的情况,此时:
应用宕机后,其持有的锁不会被释放!后续任何其他连接对于该表的任何SQL都会被持续阻塞!
VtDriver作为Vitess的客户端方案,如果其sequence实现采用事务锁的方式,由于各个应用端都会与MySQL服务直连,即各个应用获取sequence的过程都会产生锁表的行为。此时,一旦应用端由于某些原因出现锁表时长增大,甚至于应用宕机的情况,则所有应用都会由于其锁表而产生非常明显的性能下降甚至死锁。采用cas的方式使得整个过程不需要显式的开启事务,不需要锁行,自然也不存在潜在的死锁风险。当然,CAS在并发高于一定程度时会出现各个线程互相争抢资源,此时会有更新失败不断重试的情况发生,给CPU带来一定的压力,而这可以通过设置更大的cache值,增加本地缓存数量的方式来调节。