mirror of
https://github.com/dromara/hertzbeat.git
synced 2026-09-17 09:40:58 +00:00
[bugfix] redis cluster monitor master-slave relationship is inconsistent (#3874)
Co-authored-by: Tomsun28 <tomsun28@outlook.com>
This commit is contained in:
+15
-6
@@ -67,6 +67,8 @@ public class RedisCommonCollectImpl extends AbstractCollect {
|
||||
|
||||
private static final String CLUSTER = "3";
|
||||
|
||||
private static final String SINGLE = "1";
|
||||
|
||||
private static final String CLUSTER_INFO = "cluster";
|
||||
|
||||
private static final String UNIQUE_IDENTITY = "identity";
|
||||
@@ -140,7 +142,7 @@ public class RedisCommonCollectImpl extends AbstractCollect {
|
||||
* @return data
|
||||
*/
|
||||
private List<Map<String, String>> getClusterRedisInfo(Metrics metrics) throws GeneralSecurityException, IOException {
|
||||
Map<String, StatefulRedisClusterConnection<String, String>> connectionMap = getConnectionList(metrics.getRedis());
|
||||
Map<String, StatefulRedisConnection<String, String>> connectionMap = getConnectionList(metrics.getRedis());
|
||||
List<Map<String, String>> list = new ArrayList<>(connectionMap.size());
|
||||
connectionMap.forEach((identity, connection) ->{
|
||||
String info = connection.sync().info(metrics.getName());
|
||||
@@ -214,16 +216,23 @@ public class RedisCommonCollectImpl extends AbstractCollect {
|
||||
* @param redisProtocol protocol
|
||||
* @return connection map
|
||||
*/
|
||||
private Map<String, StatefulRedisClusterConnection<String, String>> getConnectionList(RedisProtocol redisProtocol) throws GeneralSecurityException, IOException {
|
||||
private Map<String, StatefulRedisConnection<String, String>> getConnectionList(RedisProtocol redisProtocol) throws GeneralSecurityException, IOException {
|
||||
// first connection
|
||||
StatefulRedisClusterConnection<String, String> connection = getClusterConnection(redisProtocol);
|
||||
Partitions partitions = connection.getPartitions();
|
||||
Map<String, StatefulRedisClusterConnection<String, String>> clusterConnectionMap = new HashMap<>(partitions.size());
|
||||
Map<String, StatefulRedisConnection<String, String>> clusterConnectionMap = new HashMap<>(partitions.size());
|
||||
for (RedisClusterNode partition : partitions) {
|
||||
RedisURI uri = partition.getUri();
|
||||
redisProtocol.setHost(uri.getHost());
|
||||
redisProtocol.setPort(String.valueOf(uri.getPort()));
|
||||
StatefulRedisClusterConnection<String, String> clusterConnection = getClusterConnection(redisProtocol);
|
||||
RedisProtocol singleRedisProtocol = RedisProtocol.builder()
|
||||
.host(uri.getHost())
|
||||
.port(String.valueOf(uri.getPort()))
|
||||
.username(redisProtocol.getUsername())
|
||||
.password(redisProtocol.getPassword())
|
||||
.pattern(SINGLE)
|
||||
.timeout(redisProtocol.getTimeout())
|
||||
.sshTunnel(redisProtocol.getSshTunnel())
|
||||
.build();
|
||||
StatefulRedisConnection<String, String> clusterConnection = getSingleConnection(singleRedisProtocol);
|
||||
clusterConnectionMap.put(doUri(uri.getHost(), uri.getPort()), clusterConnection);
|
||||
}
|
||||
return clusterConnectionMap;
|
||||
|
||||
+30
-12
@@ -19,10 +19,13 @@ package org.apache.hertzbeat.collector.collect.redis;
|
||||
|
||||
import static org.apache.hertzbeat.common.constants.CommonConstants.TYPE_STRING;
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
|
||||
import io.lettuce.core.RedisClient;
|
||||
import io.lettuce.core.RedisURI;
|
||||
import io.lettuce.core.api.StatefulRedisConnection;
|
||||
import io.lettuce.core.api.sync.RedisCommands;
|
||||
import io.lettuce.core.cluster.RedisClusterClient;
|
||||
import io.lettuce.core.cluster.api.StatefulRedisClusterConnection;
|
||||
import io.lettuce.core.cluster.api.sync.RedisAdvancedClusterCommands;
|
||||
import io.lettuce.core.cluster.models.partitions.Partitions;
|
||||
import io.lettuce.core.cluster.models.partitions.RedisClusterNode;
|
||||
import io.lettuce.core.resource.ClientResources;
|
||||
@@ -37,6 +40,7 @@ import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.extension.ExtendWith;
|
||||
import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.MockedStatic;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.junit.jupiter.MockitoExtension;
|
||||
|
||||
@@ -51,13 +55,19 @@ public class RedisClusterCollectImplTest {
|
||||
|
||||
|
||||
@Mock
|
||||
private StatefulRedisClusterConnection<String, String> connection;
|
||||
private StatefulRedisClusterConnection<String, String> clusterConnection;
|
||||
|
||||
@Mock
|
||||
private RedisAdvancedClusterCommands<String, String> cmd;
|
||||
private StatefulRedisConnection<String, String> singleConnection;
|
||||
|
||||
@Mock
|
||||
private RedisClusterClient client;
|
||||
private RedisCommands<String, String> cmd;
|
||||
|
||||
@Mock
|
||||
private RedisClusterClient clusterClient;
|
||||
|
||||
@Mock
|
||||
private RedisClient singleClient;
|
||||
|
||||
@BeforeEach
|
||||
void setUp() {
|
||||
@@ -65,8 +75,10 @@ public class RedisClusterCollectImplTest {
|
||||
|
||||
@AfterEach
|
||||
void setDown() {
|
||||
connection.close();
|
||||
client.shutdown();
|
||||
clusterConnection.close();
|
||||
singleConnection.close();
|
||||
clusterClient.shutdown();
|
||||
singleClient.shutdown();
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -110,10 +122,14 @@ public class RedisClusterCollectImplTest {
|
||||
metrics.setAliasFields(aliasField);
|
||||
metrics.setFields(fields);
|
||||
|
||||
|
||||
Mockito.mockStatic(RedisClusterClient.class).when(() -> RedisClusterClient.create(Mockito.any(ClientResources.class),
|
||||
Mockito.any(RedisURI.class))).thenReturn(client);
|
||||
Mockito.when(client.connect()).thenReturn(connection);
|
||||
MockedStatic<RedisClusterClient> redisClusterClientMockedStatic = Mockito.mockStatic(RedisClusterClient.class);
|
||||
redisClusterClientMockedStatic.when(() -> RedisClusterClient.create(Mockito.any(ClientResources.class),
|
||||
Mockito.any(RedisURI.class))).thenReturn(clusterClient);
|
||||
Mockito.when(clusterClient.connect()).thenReturn(clusterConnection);
|
||||
MockedStatic<RedisClient> redisClientMockedStatic = Mockito.mockStatic(RedisClient.class);
|
||||
redisClientMockedStatic.when(() -> RedisClient.create(Mockito.any(ClientResources.class),
|
||||
Mockito.any(RedisURI.class))).thenReturn(singleClient);
|
||||
Mockito.when(singleClient.connect()).thenReturn(singleConnection);
|
||||
|
||||
Partitions partitions = new Partitions();
|
||||
RedisClusterNode node = new RedisClusterNode();
|
||||
@@ -125,9 +141,9 @@ public class RedisClusterCollectImplTest {
|
||||
node2.setUri(RedisURI.create("redis://" + uri2));
|
||||
partitions.add(node2);
|
||||
|
||||
Mockito.when(connection.getPartitions()).thenReturn(partitions);
|
||||
Mockito.when(clusterConnection.getPartitions()).thenReturn(partitions);
|
||||
|
||||
Mockito.when(connection.sync()).thenReturn(cmd);
|
||||
Mockito.when(singleConnection.sync()).thenReturn(cmd);
|
||||
Mockito.when(cmd.info(metrics.getName())).thenReturn(info);
|
||||
Mockito.when(cmd.clusterInfo()).thenReturn(clusterInfo);
|
||||
|
||||
@@ -147,6 +163,8 @@ public class RedisClusterCollectImplTest {
|
||||
assertEquals(row.getColumns(2), uri2);
|
||||
}
|
||||
}
|
||||
redisClusterClientMockedStatic.close();
|
||||
redisClientMockedStatic.close();
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user