Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
weixin_57962205
redisson
提交
55502492
R
redisson
项目概览
weixin_57962205
/
redisson
与 Fork 源项目一致
从无法访问的项目Fork
通知
10
Star
1
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
R
redisson
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
前往新版Gitcode,体验更适合开发者的 AI 搜索 >>
提交
55502492
编写于
2月 26, 2021
作者:
N
Nikita Koksharov
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
refactoring
上级
42313bab
变更
4
隐藏空白更改
内联
并排
Showing
4 changed file
with
23 addition
and
36 deletion
+23
-36
redisson/src/main/java/org/redisson/cluster/ClusterConnectionManager.java
...n/java/org/redisson/cluster/ClusterConnectionManager.java
+3
-3
redisson/src/main/java/org/redisson/connection/MasterSlaveConnectionManager.java
...org/redisson/connection/MasterSlaveConnectionManager.java
+8
-26
redisson/src/main/java/org/redisson/connection/ReplicatedConnectionManager.java
.../org/redisson/connection/ReplicatedConnectionManager.java
+2
-2
redisson/src/main/java/org/redisson/connection/SentinelConnectionManager.java
...va/org/redisson/connection/SentinelConnectionManager.java
+10
-5
未找到文件。
redisson/src/main/java/org/redisson/cluster/ClusterConnectionManager.java
浏览文件 @
55502492
...
...
@@ -91,7 +91,7 @@ public class ClusterConnectionManager extends MasterSlaveConnectionManager {
List
<
String
>
failedMasters
=
new
ArrayList
<
String
>();
for
(
String
address
:
cfg
.
getNodeAddresses
())
{
RedisURI
addr
=
new
RedisURI
(
address
);
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
null
,
addr
.
getHost
());
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
addr
.
getHost
());
try
{
RedisConnection
connection
=
connectionFuture
.
syncUninterruptibly
().
getNow
();
...
...
@@ -276,7 +276,7 @@ public class ClusterConnectionManager extends MasterSlaveConnectionManager {
}
RPromise
<
Void
>
result
=
new
RedissonPromise
<>();
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
partition
.
getMasterAddress
(),
null
,
configEndpointHostName
);
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
partition
.
getMasterAddress
(),
configEndpointHostName
);
connectionFuture
.
onComplete
((
connection
,
ex1
)
->
{
if
(
ex1
!=
null
)
{
log
.
error
(
"Can't connect to master: {} with slot ranges: {}"
,
partition
.
getMasterAddress
(),
partition
.
getSlotRanges
());
...
...
@@ -425,7 +425,7 @@ public class ClusterConnectionManager extends MasterSlaveConnectionManager {
return
;
}
RedisURI
uri
=
iterator
.
next
();
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
uri
,
null
,
configEndpointHostName
);
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
uri
,
configEndpointHostName
);
connectionFuture
.
onComplete
((
connection
,
e
)
->
{
if
(
e
!=
null
)
{
lastException
.
set
(
e
);
...
...
redisson/src/main/java/org/redisson/connection/MasterSlaveConnectionManager.java
浏览文件 @
55502492
...
...
@@ -36,10 +36,7 @@ import org.redisson.client.codec.Codec;
import
org.redisson.client.protocol.RedisCommand
;
import
org.redisson.cluster.ClusterSlotRange
;
import
org.redisson.command.CommandSyncService
;
import
org.redisson.config.BaseMasterSlaveServersConfig
;
import
org.redisson.config.Config
;
import
org.redisson.config.MasterSlaveServersConfig
;
import
org.redisson.config.TransportMode
;
import
org.redisson.config.*
;
import
org.redisson.misc.CountableListener
;
import
org.redisson.misc.InfinitySemaphoreLatch
;
import
org.redisson.misc.RPromise
;
...
...
@@ -149,7 +146,7 @@ public class MasterSlaveConnectionManager implements ConnectionManager {
protected
PublishSubscribeService
subscribeService
;
private
final
Map
<
Object
,
RedisConnection
>
nodeConnections
=
new
ConcurrentHashMap
<>();
private
final
Map
<
RedisURI
,
RedisConnection
>
nodeConnections
=
new
ConcurrentHashMap
<>();
public
MasterSlaveConnectionManager
(
MasterSlaveServersConfig
cfg
,
Config
config
,
UUID
id
)
{
this
(
config
,
id
);
...
...
@@ -244,34 +241,19 @@ public class MasterSlaveConnectionManager implements ConnectionManager {
}
}
protected
final
void
disconnectNode
(
RedisClient
client
)
{
RedisConnection
conn
=
nodeConnections
.
remove
(
client
);
if
(
conn
!=
null
)
{
conn
.
closeAsync
();
}
}
protected
final
RFuture
<
RedisConnection
>
connectToNode
(
BaseMasterSlaveServersConfig
<?>
cfg
,
RedisURI
addr
,
RedisClient
client
,
String
sslHostname
)
{
final
Object
key
;
if
(
client
!=
null
)
{
key
=
client
;
}
else
{
key
=
addr
;
}
RedisConnection
conn
=
nodeConnections
.
get
(
key
);
protected
final
RFuture
<
RedisConnection
>
connectToNode
(
BaseConfig
<?>
cfg
,
RedisURI
addr
,
String
sslHostname
)
{
RedisConnection
conn
=
nodeConnections
.
get
(
addr
);
if
(
conn
!=
null
)
{
if
(!
conn
.
isActive
())
{
nodeConnections
.
remove
(
key
);
nodeConnections
.
remove
(
addr
);
conn
.
closeAsync
();
}
else
{
return
RedissonPromise
.
newSucceededFuture
(
conn
);
}
}
if
(
addr
!=
null
)
{
client
=
createClient
(
NodeType
.
MASTER
,
addr
,
cfg
.
getConnectTimeout
(),
cfg
.
getTimeout
(),
sslHostname
);
}
final
RPromise
<
RedisConnection
>
result
=
new
RedissonPromise
<
RedisConnection
>();
RedisClient
client
=
createClient
(
NodeType
.
MASTER
,
addr
,
cfg
.
getConnectTimeout
(),
cfg
.
getTimeout
(),
sslHostname
);
RPromise
<
RedisConnection
>
result
=
new
RedissonPromise
<>();
RFuture
<
RedisConnection
>
future
=
client
.
connectAsync
();
future
.
onComplete
((
connection
,
e
)
->
{
if
(
e
!=
null
)
{
...
...
@@ -280,7 +262,7 @@ public class MasterSlaveConnectionManager implements ConnectionManager {
}
if
(
connection
.
isActive
())
{
nodeConnections
.
put
(
key
,
connection
);
nodeConnections
.
put
(
addr
,
connection
);
result
.
trySuccess
(
connection
);
}
else
{
connection
.
closeAsync
();
...
...
redisson/src/main/java/org/redisson/connection/ReplicatedConnectionManager.java
浏览文件 @
55502492
...
...
@@ -71,7 +71,7 @@ public class ReplicatedConnectionManager extends MasterSlaveConnectionManager {
for
(
String
address
:
cfg
.
getNodeAddresses
())
{
RedisURI
addr
=
new
RedisURI
(
address
);
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
null
,
addr
.
getHost
());
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
addr
.
getHost
());
connectionFuture
.
awaitUninterruptibly
();
RedisConnection
connection
=
connectionFuture
.
getNow
();
if
(
connection
==
null
)
{
...
...
@@ -131,7 +131,7 @@ public class ReplicatedConnectionManager extends MasterSlaveConnectionManager {
AtomicInteger
count
=
new
AtomicInteger
(
cfg
.
getNodeAddresses
().
size
());
for
(
String
address
:
cfg
.
getNodeAddresses
())
{
RedisURI
addr
=
new
RedisURI
(
address
);
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
null
,
addr
.
getHost
());
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
addr
.
getHost
());
connectionFuture
.
onComplete
((
connection
,
exc
)
->
{
if
(
exc
!=
null
)
{
log
.
error
(
exc
.
getMessage
(),
exc
);
...
...
redisson/src/main/java/org/redisson/connection/SentinelConnectionManager.java
浏览文件 @
55502492
...
...
@@ -296,7 +296,7 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
}
Set
<
RedisURI
>
newUris
=
future
.
getNow
().
stream
()
.
map
(
addr
->
toURI
(
addr
.
getAddress
().
getHostAddress
(),
""
+
addr
.
getPort
()
))
.
map
(
addr
->
getIpAddr
(
addr
))
.
collect
(
Collectors
.
toSet
());
for
(
RedisURI
uri
:
newUris
)
{
...
...
@@ -343,7 +343,8 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
}
RedisClient
client
=
iterator
.
next
();
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
null
,
null
,
client
,
null
);
RedisURI
addr
=
getIpAddr
(
client
.
getAddr
());
RFuture
<
RedisConnection
>
connectionFuture
=
connectToNode
(
cfg
,
addr
,
null
);
connectionFuture
.
onComplete
((
connection
,
e
)
->
{
if
(
e
!=
null
)
{
lastException
.
set
(
e
);
...
...
@@ -484,7 +485,7 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
}).
collect
(
Collectors
.
toSet
());
InetSocketAddress
addr
=
connection
.
getRedisClient
().
getAddr
();
RedisURI
currentAddr
=
toURI
(
addr
.
getAddress
().
getHostAddress
(),
""
+
addr
.
getPort
()
);
RedisURI
currentAddr
=
getIpAddr
(
addr
);
newUris
.
add
(
currentAddr
);
updateSentinels
(
newUris
);
...
...
@@ -505,7 +506,7 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
for
(
RedisURI
uri
:
currentUris
)
{
RedisClient
sentinel
=
SentinelConnectionManager
.
this
.
sentinels
.
remove
(
uri
);
if
(
sentinel
!=
null
)
{
disconnectNode
(
sentinel
);
disconnectNode
(
uri
);
sentinel
.
shutdownAsync
();
log
.
warn
(
"sentinel: {} is down"
,
uri
);
}
...
...
@@ -540,7 +541,7 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
return
;
}
RedisURI
ipAddr
=
toURI
(
client
.
getAddr
().
getAddress
().
getHostAddress
(),
""
+
client
.
getAddr
().
getPort
());
RedisURI
ipAddr
=
getIpAddr
(
client
.
getAddr
());
if
(
isHostname
)
{
RedisClient
sentinel
=
sentinels
.
get
(
ipAddr
);
if
(
sentinel
!=
null
)
{
...
...
@@ -574,6 +575,10 @@ public class SentinelConnectionManager extends MasterSlaveConnectionManager {
return
result
;
}
private
RedisURI
getIpAddr
(
InetSocketAddress
addr
)
{
return
toURI
(
addr
.
getAddress
().
getHostAddress
(),
""
+
addr
.
getPort
());
}
private
RFuture
<
Void
>
addSlave
(
RedisURI
uri
)
{
RPromise
<
Void
>
result
=
new
RedissonPromise
<
Void
>();
// to avoid addition twice
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录