-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathclient_zk.go
More file actions
207 lines (186 loc) · 4.5 KB
/
client_zk.go
File metadata and controls
207 lines (186 loc) · 4.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
package distlock
import (
"bytes"
"time"
"github.com/rs/zerolog/log"
"github.com/samuel/go-zookeeper/zk"
)
const (
_LockerRootPath = "/distlock"
_LockerLockPathFastLockUsed = "/distlock/fast-lock"
_LockerLockPathFastLockUsedPrefix = "/distlock/fast-lock/request-"
_LockerLockPathFastLockUsedShortestPrefix = "request-"
_MaxRetries = 3
)
// EstablishZKConn 建立一条连接zookeeper集群的TCP连接.
func EstablishZKConn(endpoints []string) (*zk.Conn, <-chan zk.Event) {
conn, evCh, err := zk.Connect(endpoints, time.Second*time.Duration(10))
if err != nil {
log.Fatal().Err(err).Msgf("failed to connect to zookeeper cluster (%v)", endpoints)
}
createIfNotExistOrDie(conn, _LockerRootPath)
createIfNotExistOrDie(conn, _LockerLockPathFastLockUsed)
return conn, evCh
}
// CloseZKConn 关闭TCP连接.
func CloseZKConn(conn *zk.Conn) {
conn.Close()
}
// 如果ZNode不存在就创建
func createIfNotExist(conn *zk.Conn, path string) error {
if _, err := safeCreate(conn, path, []byte(""), 0); err != nil && err != zk.ErrNodeExists {
return err
}
return nil
}
// 如果ZNode不存在就创建, 出错直接panic
func createIfNotExistOrDie(conn *zk.Conn, path string) {
if err := createIfNotExist(conn, path); err != nil {
log.Fatal().Err(err).Msgf("failed to create znode <%s>", path)
}
}
// 创建ZNode
func safeCreate(conn *zk.Conn, path string, data []byte, flags int32) (string, error) {
var _path string
var err error
var _data []byte
var retry bool
LOOP:
for i := 0; i < _MaxRetries; i++ {
_path, err = conn.Create(path, data, flags, zk.WorldACL(zk.PermAll))
switch err {
// No need to search for the node since it can't exist. Just try again.
case zk.ErrSessionExpired:
{
continue
}
// 连接关闭, 可能因为暂时的网络问题, 直接重试
case zk.ErrConnectionClosed:
{
retry = true
continue
}
// ZNode已存在
case zk.ErrNodeExists:
{
// 之前就创建过
if !retry {
return _path, zk.ErrNodeExists
}
// 因为网络问题导致的假失败
_data, _, err = safeGet(conn, path)
if err != nil {
// 又可能因为暂时的网络问题, 请重试
continue
}
if bytes.Equal(data, _data) {
return _path, nil
}
return "", zk.ErrUnknown
}
// TODO: 处理更多的错误情形
default:
{
break LOOP
}
}
}
return _path, err
}
// 获取ZNode的值
func safeGet(conn *zk.Conn, path string) ([]byte, *zk.Stat, error) {
var _data []byte
var _stat *zk.Stat
var err error
LOOP:
for i := 0; i < _MaxRetries; i++ {
_data, _stat, err = conn.Get(path)
switch err {
// session过期直接panic
case zk.ErrSessionExpired:
{
log.Fatal().Err(zk.ErrSessionExpired).Msgf("failed to get value of znode <%s>", path)
}
// 连接关闭, 可能因为暂时的网络问题, 请重试
case zk.ErrConnectionClosed:
{
continue
}
// TODO: 处理更多的错误情形
default:
{
break LOOP
}
}
}
return _data, _stat, err
}
// 删除ZNode
func safeDelete(conn *zk.Conn, path string, version int32) error {
var err error
var retry bool
LOOP:
for i := 0; i < _MaxRetries; i++ {
err = conn.Delete(path, version)
switch err {
// session过期直接panic
case zk.ErrSessionExpired:
{
log.Fatal().Err(zk.ErrSessionExpired).Msgf("failed to delete znode <%s>", path)
}
// 连接关闭, 可能因为暂时的网络问题, 请重试
case zk.ErrConnectionClosed:
{
retry = true
continue
}
// ZNode不存在
case zk.ErrNoNode:
{
// 因为网络问题导致的假失败
if retry {
return nil
}
return zk.ErrNoNode
}
// TODO: 处理更多的错误情形
default:
{
break LOOP
}
}
}
return err
}
// 获取ZNode所有的子节点
func safeGetChildren(conn *zk.Conn, path string, watch bool) ([]string, <-chan zk.Event, error) {
var _children []string
var _watcher <-chan zk.Event
var err error
LOOP:
for i := 0; i < _MaxRetries; i++ {
if watch {
_children, _, _watcher, err = conn.ChildrenW(path)
} else {
_children, _, err = conn.Children(path)
}
switch err {
// session过期直接panic
case zk.ErrSessionExpired:
{
log.Fatal().Err(zk.ErrSessionExpired).Msgf("failed to list children of znode <%s>", path)
}
// 连接关闭, 可能因为暂时的网络问题, 请重试
case zk.ErrConnectionClosed:
{
continue
}
// TODO: 处理更多的错误情形
default:
{
break LOOP
}
}
}
return _children, _watcher, err
}