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
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
| import java.io.*;
import java.net.*;
import java.util.*;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* Gossip节点基础实现
*/
public class GossipNode {
private final String nodeId;
private final String address;
private final int port;
private final int gossipInterval; // 毫秒
private final int failureTimeout; // 毫秒
private final int fanout; // 每次gossip的目标节点数
// 节点状态管理
private final Map<String, NodeState> memberList = new ConcurrentHashMap<>();
private final Map<String, GossipMessage> messageStore = new ConcurrentHashMap<>();
private final Set<String> processedMessages = ConcurrentHashMap.newKeySet();
// 线程管理
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(4);
private final ExecutorService executor = Executors.newCachedThreadPool();
private final AtomicBoolean running = new AtomicBoolean(false);
private ServerSocket serverSocket;
public GossipNode(String nodeId, String address, int port) {
this.nodeId = nodeId;
this.address = address;
this.port = port;
this.gossipInterval = 1000; // 1秒
this.failureTimeout = 5000; // 5秒
this.fanout = 3; // 每次选择3个节点
// 添加自己到成员列表
memberList.put(nodeId, new NodeState(nodeId, address, port));
}
/**
* 启动Gossip节点
*/
public void start() throws IOException {
if (running.compareAndSet(false, true)) {
// 启动服务器监听
serverSocket = new ServerSocket(port);
executor.submit(this::acceptConnections);
// 启动定时任务
scheduler.scheduleAtFixedRate(this::gossipRound,
gossipInterval, gossipInterval, TimeUnit.MILLISECONDS);
scheduler.scheduleAtFixedRate(this::failureDetection,
failureTimeout, failureTimeout, TimeUnit.MILLISECONDS);
scheduler.scheduleAtFixedRate(this::cleanup,
30000, 30000, TimeUnit.MILLISECONDS);
System.out.println("Gossip节点启动: " + nodeId + " 在 " + address + ":" + port);
}
}
/**
* 停止Gossip节点
*/
public void stop() throws IOException {
if (running.compareAndSet(true, false)) {
scheduler.shutdown();
executor.shutdown();
if (serverSocket != null) {
serverSocket.close();
}
System.out.println("Gossip节点停止: " + nodeId);
}
}
/**
* 接受连接
*/
private void acceptConnections() {
while (running.get() && !serverSocket.isClosed()) {
try {
Socket clientSocket = serverSocket.accept();
executor.submit(() -> handleConnection(clientSocket));
} catch (IOException e) {
if (running.get()) {
System.err.println("接受连接错误: " + e.getMessage());
}
}
}
}
/**
* 处理连接
*/
private void handleConnection(Socket socket) {
try (ObjectInputStream input = new ObjectInputStream(socket.getInputStream());
ObjectOutputStream output = new ObjectOutputStream(socket.getOutputStream())) {
Object request = input.readObject();
if (request instanceof GossipMessage) {
handleGossipMessage((GossipMessage) request);
output.writeObject("ACK");
} else if (request instanceof MembershipRequest) {
handleMembershipRequest((MembershipRequest) request, output);
}
} catch (Exception e) {
System.err.println("处理连接错误: " + e.getMessage());
} finally {
try {
socket.close();
} catch (IOException e) {
// 忽略关闭错误
}
}
}
/**
* Gossip轮次
*/
private void gossipRound() {
try {
// 更新自己的心跳
memberList.put(nodeId, memberList.get(nodeId).updateHeartbeat());
// 选择随机节点进行gossip
List<String> targetNodes = selectRandomNodes(fanout);
for (String targetNodeId : targetNodes) {
NodeState targetNode = memberList.get(targetNodeId);
if (targetNode != null && !targetNode.getNodeId().equals(nodeId)) {
executor.submit(() -> gossipToNode(targetNode));
}
}
} catch (Exception e) {
System.err.println("Gossip轮次错误: " + e.getMessage());
}
}
/**
* 向指定节点发送gossip
*/
private void gossipToNode(NodeState targetNode) {
try (Socket socket = new Socket(targetNode.getAddress(), targetNode.getPort());
ObjectOutputStream output = new ObjectOutputStream(socket.getOutputStream());
ObjectInputStream input = new ObjectInputStream(socket.getInputStream())) {
// 发送成员列表更新
MembershipRequest request = new MembershipRequest(nodeId, new ArrayList<>(memberList.values()));
output.writeObject(request);
Object response = input.readObject();
if (response instanceof MembershipResponse) {
handleMembershipResponse((MembershipResponse) response);
}
} catch (Exception e) {
// 标记节点为可疑
System.err.println("与节点 " + targetNode.getNodeId() + " 通信失败: " + e.getMessage());
}
}
/**
* 选择随机节点
*/
private List<String> selectRandomNodes(int count) {
List<String> allNodes = new ArrayList<>(memberList.keySet());
allNodes.remove(nodeId); // 移除自己
Collections.shuffle(allNodes);
return allNodes.subList(0, Math.min(count, allNodes.size()));
}
/**
* 处理Gossip消息
*/
private void handleGossipMessage(GossipMessage message) {
if (processedMessages.contains(message.getMessageId()) || message.getTtl() <= 0) {
return;
}
processedMessages.add(message.getMessageId());
messageStore.put(message.getMessageId(), message);
// 处理消息内容
processMessage(message);
// 继续传播消息
if (message.getTtl() > 1) {
propagateMessage(message.decrementTTL());
}
}
/**
* 传播消息
*/
private void propagateMessage(GossipMessage message) {
List<String> targetNodes = selectRandomNodes(fanout);
for (String targetNodeId : targetNodes) {
NodeState targetNode = memberList.get(targetNodeId);
if (targetNode != null && !targetNode.getNodeId().equals(nodeId)) {
executor.submit(() -> sendMessage(targetNode, message));
}
}
}
/**
* 发送消息到节点
*/
private void sendMessage(NodeState targetNode, GossipMessage message) {
try (Socket socket = new Socket(targetNode.getAddress(), targetNode.getPort());
ObjectOutputStream output = new ObjectOutputStream(socket.getOutputStream());
ObjectInputStream input = new ObjectInputStream(socket.getInputStream())) {
output.writeObject(message);
input.readObject(); // 读取ACK
} catch (Exception e) {
System.err.println("发送消息到节点 " + targetNode.getNodeId() + " 失败: " + e.getMessage());
}
}
/**
* 处理成员关系请求
*/
private void handleMembershipRequest(MembershipRequest request, ObjectOutputStream output) throws IOException {
// 合并成员列表
for (NodeState nodeState : request.getMemberList()) {
NodeState existing = memberList.get(nodeState.getNodeId());
if (existing == null || nodeState.getLastHeartbeat() > existing.getLastHeartbeat()) {
memberList.put(nodeState.getNodeId(), nodeState);
}
}
// 返回自己的成员列表
MembershipResponse response = new MembershipResponse(nodeId, new ArrayList<>(memberList.values()));
output.writeObject(response);
}
/**
* 处理成员关系响应
*/
private void handleMembershipResponse(MembershipResponse response) {
// 合并成员列表
for (NodeState nodeState : response.getMemberList()) {
NodeState existing = memberList.get(nodeState.getNodeId());
if (existing == null || nodeState.getLastHeartbeat() > existing.getLastHeartbeat()) {
memberList.put(nodeState.getNodeId(), nodeState);
}
}
}
/**
* 故障检测
*/
private void failureDetection() {
long currentTime = System.currentTimeMillis();
List<String> failedNodes = new ArrayList<>();
for (Map.Entry<String, NodeState> entry : memberList.entrySet()) {
String nodeId = entry.getKey();
NodeState nodeState = entry.getValue();
if (!nodeId.equals(this.nodeId) &&
currentTime - nodeState.getLastHeartbeat() > failureTimeout) {
failedNodes.add(nodeId);
}
}
// 移除失效节点
for (String failedNode : failedNodes) {
memberList.remove(failedNode);
System.out.println("检测到节点失效: " + failedNode);
}
}
/**
* 清理过期消息
*/
private void cleanup() {
long currentTime = System.currentTimeMillis();
long maxAge = 300000; // 5分钟
messageStore.entrySet().removeIf(entry ->
currentTime - entry.getValue().getTimestamp() > maxAge);
// 清理处理过的消息ID(保留最近的1000个)
if (processedMessages.size() > 1000) {
List<String> messageIds = new ArrayList<>(processedMessages);
processedMessages.clear();
processedMessages.addAll(messageIds.subList(messageIds.size() - 500, messageIds.size()));
}
}
/**
* 处理消息内容
*/
private void processMessage(GossipMessage message) {
System.out.println("收到消息: " + message.getContent() + " 来自: " + message.getSourceNode());
// 在这里可以添加具体的业务逻辑
// 例如:状态更新、事件处理等
}
/**
* 广播消息
*/
public void broadcast(String content) {
String messageId = UUID.randomUUID().toString();
GossipMessage message = new GossipMessage(messageId, content, nodeId, 10);
processedMessages.add(messageId);
messageStore.put(messageId, message);
propagateMessage(message);
System.out.println("广播消息: " + content);
}
/**
* 加入集群
*/
public void joinCluster(String seedAddress, int seedPort) throws IOException, ClassNotFoundException {
try (Socket socket = new Socket(seedAddress, seedPort);
ObjectOutputStream output = new ObjectOutputStream(socket.getOutputStream());
ObjectInputStream input = new ObjectInputStream(socket.getInputStream())) {
MembershipRequest request = new MembershipRequest(nodeId, Arrays.asList(memberList.get(nodeId)));
output.writeObject(request);
Object response = input.readObject();
if (response instanceof MembershipResponse) {
handleMembershipResponse((MembershipResponse) response);
System.out.println("成功加入集群,当前成员数: " + memberList.size());
}
}
}
/**
* 获取集群状态
*/
public void printClusterStatus() {
System.out.println("\n=== 集群状态 ===");
System.out.println("节点ID: " + nodeId);
System.out.println("成员数量: " + memberList.size());
System.out.println("消息数量: " + messageStore.size());
System.out.println("\n成员列表:");
for (NodeState node : memberList.values()) {
long timeSinceHeartbeat = System.currentTimeMillis() - node.getLastHeartbeat();
System.out.printf("- %s (%s:%d) 心跳: %dms前\n",
node.getNodeId(), node.getAddress(), node.getPort(), timeSinceHeartbeat);
}
System.out.println("================\n");
}
// Getters
public String getNodeId() { return nodeId; }
public int getMemberCount() { return memberList.size(); }
public boolean isRunning() { return running.get(); }
}
/**
* 成员关系请求
*/
class MembershipRequest implements Serializable {
private final String requestorId;
private final List<NodeState> memberList;
public MembershipRequest(String requestorId, List<NodeState> memberList) {
this.requestorId = requestorId;
this.memberList = memberList;
}
public String getRequestorId() { return requestorId; }
public List<NodeState> getMemberList() { return memberList; }
}
/**
* 成员关系响应
*/
class MembershipResponse implements Serializable {
private final String responderId;
private final List<NodeState> memberList;
public MembershipResponse(String responderId, List<NodeState> memberList) {
this.responderId = responderId;
this.memberList = memberList;
}
public String getResponderId() { return responderId; }
public List<NodeState> getMemberList() { return memberList; }
}
|