Browse Source

Merge pull request #33 from L-fushen/develop-go-regworkerid

增加redis哨兵和集群模式的支持
pull/19/MERGE
yitter GitHub 2 years ago
parent
commit
acc99529e1
No known key found for this signature in database GPG Key ID: 4AEE18F83AFDEB23
2 changed files with 100 additions and 74 deletions
  1. +43
    -21
      Go/regworkerid/main.go
  2. +57
    -53
      Go/regworkerid/regworkerid/reghelper.go

+ 43
- 21
Go/regworkerid/main.go View File

@@ -9,53 +9,75 @@ import (
) )


func main() { func main() {
ip := "localhost"

address := "127.0.0.1:6379"
password := "" password := ""
masterName := ""


ipChar := C.CString(ip)
addrChar := C.CString(address)
passChar := C.CString(password) passChar := C.CString(password)
masterNameChar := C.CString(masterName)


workerIdList := RegisterMany(ipChar, 6379, passChar, 4, 3, 0)
workerIdList := RegisterMany(addrChar, passChar, 4, masterNameChar, 3, 5)
for _, value := range workerIdList { for _, value := range workerIdList {
fmt.Println("注册的WorkerId:", value) fmt.Println("注册的WorkerId:", value)
} }


id := RegisterOne(ipChar, 6379, passChar, 4, 0)
id := RegisterOne(addrChar, passChar, 4, masterNameChar, 0)
fmt.Println("注册的WorkerId:", id) fmt.Println("注册的WorkerId:", id)


// var workerId = regworkerid.RegisterOne(ip, 6379, "", 4)
// fmt.Println("注册的WorkerId:", workerId)

fmt.Println("end") fmt.Println("end")
time.Sleep(time.Duration(300) * time.Second) time.Sleep(time.Duration(300) * time.Second)
} }


// RegisterOne 注册一个 WorkerId,会先注销所有本机已注册的记录
// address: Redis连接地址,单机模式示例:127.0.0.1:6379,哨兵/集群模式示例:127.0.0.1:26380,127.0.0.1:26381,127.0.0.1:26382
// password: Redis连接密码
// db: Redis指定存储库,示例:1
// sentinelMasterName: Redis 哨兵模式下的服务名称,示例:mymaster,非哨兵模式传入空字符串即可
// maxWorkerId: WorkerId 最大值,示例:63
//export RegisterOne //export RegisterOne
// 注册一个 WorkerId
func RegisterOne(ip *C.char, port int32, password *C.char, maxWorkerId int32, database int) int32 {
return regworkerid.RegisterOne(C.GoString(ip), port, C.GoString(password), maxWorkerId, database)
func RegisterOne(address *C.char, password *C.char, db int, sentinelMasterName *C.char, maxWorkerId int32) int32 {
return regworkerid.RegisterOne(regworkerid.RegisterConf{
Address: C.GoString(address),
Password: C.GoString(password),
DB: db,
MasterName: C.GoString(sentinelMasterName),
MaxWorkerId: maxWorkerId,
})
}

// RegisterMany 注册多个 WorkerId,会先注销所有本机已注册的记录
// address: Redis连接地址,单机模式示例:127.0.0.1:6379,哨兵/集群模式示例:127.0.0.1:26380,127.0.0.1:26381,127.0.0.1:26382
// password: Redis连接密码
// db: Redis指定存储库,示例:1
// sentinelMasterName: Redis 哨兵模式下的服务名称,示例:mymaster,非哨兵模式传入空字符串即可
// maxWorkerId: WorkerId 最大值,示例:63
// totalCount: 获取N个WorkerId,示例:5
//export RegisterMany
func RegisterMany(address *C.char, password *C.char, db int, sentinelMasterName *C.char, maxWorkerId int32, totalCount int32) []int32 {
return regworkerid.RegisterMany(regworkerid.RegisterConf{
Address: C.GoString(address),
Password: C.GoString(password),
DB: db,
MasterName: C.GoString(sentinelMasterName),
MaxWorkerId: maxWorkerId,
TotalCount: totalCount,
})
} }


// UnRegister 注销本机已注册的 WorkerId
//export UnRegister //export UnRegister
// 注销本机已注册的 WorkerId
func UnRegister() { func UnRegister() {
regworkerid.UnRegister() regworkerid.UnRegister()
} }


// export Validate
// 检查本地WorkerId是否有效(0-有效,其它-无效)
// Validate 检查本地WorkerId是否有效(0-有效,其它-无效)
//export Validate
func Validate(workerId int32) int32 { func Validate(workerId int32) int32 {
return regworkerid.Validate(workerId) return regworkerid.Validate(workerId)
} }


// RegisterMany
// 注册多个 WorkerId,会先注销所有本机已注册的记录
func RegisterMany(ip *C.char, port int32, password *C.char, maxWorkerId, totalCount int32, database int) []int32 {
// return (*C.int)(unsafe.Pointer(&values))
//return regworkerid.RegisterMany(ip, port, password, maxWorkerId, totalCount, database)
return regworkerid.RegisterMany(C.GoString(ip), port, C.GoString(password), maxWorkerId, totalCount, database)
}

// To Build a dll/so: // To Build a dll/so:


// windows: // windows:


+ 57
- 53
Go/regworkerid/regworkerid/reghelper.go View File

@@ -1,22 +1,16 @@
/*
* 版权属于:yitter(yitter@126.com)
* 开源地址:https://github.com/yitter/idgenerator
*/

// Package regworkerid implements a simple distributed id generator.
package regworkerid package regworkerid


import ( import (
"context" "context"
"fmt" "fmt"
"github.com/go-redis/redis/v8"
"strconv" "strconv"
"strings"
"sync" "sync"
"time" "time"

"github.com/go-redis/redis/v8"
) )


var _client *redis.Client
var _client redis.UniversalClient
var _ctx = context.Background() var _ctx = context.Background()
var _workerIdLock sync.Mutex var _workerIdLock sync.Mutex


@@ -29,18 +23,26 @@ var _WorkerIdLifeTimeSeconds = 15 // IdGen:WorkerId:Value:xx 的值在 redis
var _MaxLoopCount = 10 // 最大循环次数(无可用WorkerId时循环查找) var _MaxLoopCount = 10 // 最大循环次数(无可用WorkerId时循环查找)
var _SleepMillisecondEveryLoop = 200 // 每次循环后,暂停时间 var _SleepMillisecondEveryLoop = 200 // 每次循环后,暂停时间
var _MaxWorkerId int32 = 0 // 最大WorkerId值,超过此值从0开始 var _MaxWorkerId int32 = 0 // 最大WorkerId值,超过此值从0开始
var _Database int = 0 // 最大WorkerId值,超过此值从0开始


var _RedisConnString = "" var _RedisConnString = ""
var _RedisPassword = "" var _RedisPassword = ""
var _RedisDB = 0
var _RedisMasterName = ""


const _WorkerIdIndexKey string = "IdGen:WorkerId:Index" // redis 中的key const _WorkerIdIndexKey string = "IdGen:WorkerId:Index" // redis 中的key
const _WorkerIdValueKeyPrefix string = "IdGen:WorkerId:Value:" // redis 中的key const _WorkerIdValueKeyPrefix string = "IdGen:WorkerId:Value:" // redis 中的key
const _WorkerIdFlag = "Y" // IdGen:WorkerId:Value:xx 的值(将来可用 _token 替代) const _WorkerIdFlag = "Y" // IdGen:WorkerId:Value:xx 的值(将来可用 _token 替代)
const _Log = false // 是否输出日志 const _Log = false // 是否输出日志


// export Validate
// 检查本地WorkerId是否有效(0-有效,其它-无效)
type RegisterConf struct {
Address string // 注意:哨兵模式下,这里传入的是 Sentinel 节点,不是 Redis 节点
Password string
DB int
MasterName string // 注意:哨兵模式下,这里必须传入 Sentinel 服务名称
MaxWorkerId int32
TotalCount int32 // 注意:仅对 RegisterMany 生效
}

func Validate(workerId int32) int32 { func Validate(workerId int32) int32 {
for _, value := range _workerIdList { for _, value := range _workerIdList {
if value == workerId { if value == workerId {
@@ -50,15 +52,13 @@ func Validate(workerId int32) int32 {


return 0 return 0


// if workerId == _usingWorkerId {
//if workerId == _usingWorkerId {
// return 0 // return 0
// } else {
//} else {
// return -1 // return -1
// }
//}
} }


// export UnRegister
// 注销本机已注册的 WorkerId
func UnRegister() { func UnRegister() {
_workerIdLock.Lock() _workerIdLock.Lock()


@@ -80,21 +80,24 @@ func autoUnRegister() {
} }
} }


func RegisterMany(ip string, port int32, password string, maxWorkerId int32, totalCount int32, database int) []int32 {
if maxWorkerId < 0 {
func RegisterMany(conf RegisterConf) []int32 {
if conf.MaxWorkerId < 0 {
return []int32{-2} return []int32{-2}
} }


if totalCount < 1 {
if conf.TotalCount < 1 {
return []int32{-1} return []int32{-1}
} else if conf.TotalCount == 0 {
conf.TotalCount = 1
} }


autoUnRegister() autoUnRegister()


_MaxWorkerId = maxWorkerId
_RedisConnString = ip + ":" + strconv.Itoa(int(port))
_RedisPassword = password
_Database = database
_MaxWorkerId = conf.MaxWorkerId
_RedisConnString = conf.Address
_RedisPassword = conf.Password
_RedisDB = conf.DB
_RedisMasterName = conf.MasterName
_client = newRedisClient() _client = newRedisClient()
if _client == nil { if _client == nil {
return []int32{-1} return []int32{-1}
@@ -104,18 +107,18 @@ func RegisterMany(ip string, port int32, password string, maxWorkerId int32, tot
_ = _client.Close() _ = _client.Close()
} }
}() }()
// _, err := _client.Ping(_ctx).Result()
// if err != nil {
//_, err := _client.Ping(_ctx).Result()
//if err != nil {
// //panic("init redis error") // //panic("init redis error")
// return []int{-3} // return []int{-3}
// } else {
//} else {
// if _Log { // if _Log {
// fmt.Println("init redis ok") // fmt.Println("init redis ok")
// } // }
// }
//}


_lifeIndex++ _lifeIndex++
_workerIdList = make([]int32, totalCount)
_workerIdList = make([]int32, conf.TotalCount)
for key := range _workerIdList { for key := range _workerIdList {
_workerIdList[key] = -1 // 全部初始化-1 _workerIdList[key] = -1 // 全部初始化-1
} }
@@ -125,7 +128,7 @@ func RegisterMany(ip string, port int32, password string, maxWorkerId int32, tot
id := register(_lifeIndex) id := register(_lifeIndex)
if id > -1 { if id > -1 {
useExtendFunc = true useExtendFunc = true
_workerIdList[key] = id // = append(_workerIdList, id)
_workerIdList[key] = id //= append(_workerIdList, id)
} else { } else {
break break
} }
@@ -138,20 +141,19 @@ func RegisterMany(ip string, port int32, password string, maxWorkerId int32, tot
return _workerIdList return _workerIdList
} }


// export RegisterOne
// 注册一个 WorkerId,会先注销所有本机已注册的记录
func RegisterOne(ip string, port int32, password string, maxWorkerId int32, database int) int32 {
if maxWorkerId < 0 {
func RegisterOne(conf RegisterConf) int32 {
if conf.MaxWorkerId < 0 {
return -2 return -2
} }


autoUnRegister() autoUnRegister()


_MaxWorkerId = maxWorkerId
_RedisConnString = ip + ":" + strconv.Itoa(int(port))
_RedisPassword = password
_MaxWorkerId = conf.MaxWorkerId
_RedisConnString = conf.Address
_RedisPassword = conf.Password
_RedisDB = conf.DB
_RedisMasterName = conf.MasterName
_loopCount = 0 _loopCount = 0
_Database = database
_client = newRedisClient() _client = newRedisClient()
if _client == nil { if _client == nil {
return -3 return -3
@@ -161,15 +163,15 @@ func RegisterOne(ip string, port int32, password string, maxWorkerId int32, data
_ = _client.Close() _ = _client.Close()
} }
}() }()
// _, err := _client.Ping(_ctx).Result()
// if err != nil {
//_, err := _client.Ping(_ctx).Result()
//if err != nil {
// // panic("init redis error") // // panic("init redis error")
// return -3 // return -3
// } else {
//} else {
// if _Log { // if _Log {
// fmt.Println("init redis ok") // fmt.Println("init redis ok")
// } // }
// }
//}


_lifeIndex++ _lifeIndex++
var id = register(_lifeIndex) var id = register(_lifeIndex)
@@ -186,16 +188,18 @@ func register(lifeTime int) int32 {
return getNextWorkerId(lifeTime) return getNextWorkerId(lifeTime)
} }


func newRedisClient() *redis.Client {
return redis.NewClient(&redis.Options{
Addr: _RedisConnString,
Password: _RedisPassword,
DB: _Database,
// PoolSize: 1000,
// ReadTimeout: time.Millisecond * time.Duration(100),
// WriteTimeout: time.Millisecond * time.Duration(100),
// IdleTimeout: time.Second * time.Duration(60),
func newRedisClient() redis.UniversalClient {
client := redis.NewUniversalClient(&redis.UniversalOptions{
Addrs: strings.Split(_RedisConnString, ","),
Password: _RedisPassword,
DB: _RedisDB,
MasterName: _RedisMasterName,
//PoolSize: 1000,
//ReadTimeout: time.Millisecond * time.Duration(100),
//WriteTimeout: time.Millisecond * time.Duration(100),
//IdleTimeout: time.Second * time.Duration(60),
}) })
return client
} }


func getNextWorkerId(lifeTime int) int32 { func getNextWorkerId(lifeTime int) int32 {
@@ -321,9 +325,9 @@ func extendWorkerIdLifeTime(lifeIndex int, workerId int32) {
} }


// 已经被注销,则终止(此步是上一步的二次验证) // 已经被注销,则终止(此步是上一步的二次验证)
// if _usingWorkerId < 0 {
//if _usingWorkerId < 0 {
// break // break
// }
//}


// 延长 redis 数据有效期 // 延长 redis 数据有效期
extendWorkerIdFlag(myWorkerId) extendWorkerIdFlag(myWorkerId)


Loading…
Cancel
Save