尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

第1讲:分布式基石——网络通信与节点发现

第1讲:分布式基石——网络通信与节点发现 欢迎来到 MiniKV 系列第一讲我们要打好分布式系统的基础。单机程序里函数调用就是通信但在分布式系统中节点之间通过网络传递消息——这带来了全新的挑战网络延迟、消息丢失、节点宕机。这一讲我们将实现 MiniKV 的网络层和节点发现机制让多个进程能互相发现并通信。一、整体架构┌─────────────────────────────────────────────┐ │ MiniKV 节点 │ │ │ │ ┌─────────┐ ┌──────────┐ ┌───────────┐ │ │ │ Node │ │ Router │ │ Discovery │ │ │ │ Service │◄─┤ Layer │◄─┤ Service │ │ │ └────┬────┘ └────┬─────┘ └─────┬─────┘ │ │ │ │ │ │ │ ┌────▼────────────▼───────────────▼─────┐ │ │ │ Transport Layer │ │ │ │ (TCP/UDP Message Serialization) │ │ │ └───────────────────────────────────────┘ │ └─────────────────────────────────────────────┘二、消息定义与序列化2.1 消息类型# minikv/transport/message.py from enum import Enum, auto from dataclasses import dataclass, field from typing import Any, Dict, Optional import json import struct import time import uuid class MessageType(Enum): # 节点发现 PING auto() PONG auto() JOIN auto() JOIN_RESPONSE auto() LEAVE auto() # Raft 共识 REQUEST_VOTE auto() VOTE_RESPONSE auto() APPEND_ENTRIES auto() APPEND_RESPONSE auto() # KV 操作 GET auto() GET_RESPONSE auto() PUT auto() PUT_RESPONSE auto() DELETE auto() DELETE_RESPONSE auto() # 集群管理 TRANSFER_LEADER auto() ADD_MEMBER auto() REMOVE_MEMBER auto() dataclass class Message: 网络消息 msg_id: str msg_type: MessageType MessageType.PING sender_id: str receiver_id: str timestamp: float 0.0 body: Dict[str, Any] field(default_factorydict) def __post_init__(self): if not self.msg_id: self.msg_id str(uuid.uuid4()) if not self.timestamp: self.timestamp time.time() def serialize(self) - bytes: 序列化为字节流 data { msg_id: self.msg_id, msg_type: self.msg_type.name, sender_id: self.sender_id, receiver_id: self.receiver_id, timestamp: self.timestamp, body: self.body, } return json.dumps(data).encode(utf-8) staticmethod def deserialize(data: bytes) - Message: 从字节流反序列化 decoded json.loads(data.decode(utf-8)) return Message( msg_iddecoded[msg_id], msg_typeMessageType[decoded[msg_type]], sender_iddecoded[sender_id], receiver_iddecoded[receiver_id], timestampdecoded[timestamp], bodydecoded.get(body, {}) ) def reply(self, body: Dict None) - Message: 生成回复消息 reply_body body or {} reply_body[reply_to] self.msg_id return Message( msg_typeself._get_reply_type(), sender_idself.receiver_id, receiver_idself.sender_id, bodyreply_body ) def _get_reply_type(self) - MessageType: 根据请求类型获取回复类型 reply_map { MessageType.PING: MessageType.PONG, MessageType.JOIN: MessageType.JOIN_RESPONSE, MessageType.REQUEST_VOTE: MessageType.VOTE_RESPONSE, MessageType.APPEND_ENTRIES: MessageType.APPEND_RESPONSE, MessageType.GET: MessageType.GET_RESPONSE, MessageType.PUT: MessageType.PUT_RESPONSE, MessageType.DELETE: MessageType.DELETE_RESPONSE, } return reply_map.get(self.msg_type, MessageType.PONG)2.2 二进制协议高性能版本# minikv/transport/binary_protocol.py import struct import json from typing import Dict, Any class BinaryProtocol: 二进制协议可选的高性能版本 格式 ┌────────┬──────────┬──────────┬──────────┐ │ Magic │ Version │ MsgType │ BodyLen │ │ (4B) │ (1B) │ (2B) │ (4B) │ ├────────┴──────────┴──────────┴──────────┤ │ Body (JSON) │ │ BodyLen bytes │ └──────────────────────────────────────────┘ MAGIC bMKV1 HEADER_SIZE 11 staticmethod def encode(msg_type: int, body: Dict[str, Any]) - bytes: 编码消息 body_bytes json.dumps(body).encode(utf-8) header struct.pack(!4sBH, BinaryProtocol.MAGIC, 1, # version msg_type, len(body_bytes) ) return header body_bytes staticmethod def decode(data: bytes) - tuple: 解码消息返回 (msg_type, body) if len(data) BinaryProtocol.HEADER_SIZE: raise ValueError(数据太短) magic data[:4] if magic ! BinaryProtocol.MAGIC: raise ValueError(fMagic number 不匹配: {magic}) version data[4] msg_type struct.unpack_from(!H, data, 5)[0] body_len struct.unpack_from(!I, data, 7)[0] if len(data) BinaryProtocol.HEADER_SIZE body_len: raise ValueError(Body 数据不完整) body json.loads(data[BinaryProtocol.HEADER_SIZE:BinaryProtocol.HEADER_SIZE body_len]) return msg_type, body三、传输层实现3.1 TCP 连接管理# minikv/transport/tcp_transport.py import socket import threading import select import queue import logging from typing import Dict, Callable, Optional from .message import Message, MessageType logger logging.getLogger(__name__) class Connection: TCP 连接封装 def __init__(self, sock: socket.socket, addr: tuple, node_id: str ): self.sock sock self.addr addr self.node_id node_id self.send_queue queue.Queue() self.recv_buffer b self.closed False self.last_active time.time() def send(self, msg: Message): 发送消息放入队列 if not self.closed: self.send_queue.put(msg) def send_sync(self, msg: Message): 同步发送 try: data msg.serialize() # 4字节长度前缀 数据 packet struct.pack(!I, len(data)) data self.sock.sendall(packet) self.last_active time.time() except Exception as e: logger.error(f发送失败: {e}) self.close() def close(self): 关闭连接 self.closed True try: self.sock.close() except: pass class TCPServer: TCP 服务器 def __init__(self, host: str 0.0.0.0, port: int 9000): self.host host self.port port self.server_sock: Optional[socket.socket] None self.connections: Dict[str, Connection] {} # node_id - Connection self.message_handler: Optional[Callable] None self.running False self.lock threading.Lock() def start(self): 启动服务器 self.server_sock socket.socket(socket.AF_INET, socket.SOCK_STREAM) self.server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) self.server_sock.bind((self.host, self.port)) self.server_sock.listen(128) self.server_sock.setblocking(False) self.running True logger.info(fTCP Server listening on {self.host}:{self.port}) # 启动 acceptor 线程 accept_thread threading.Thread(targetself._accept_loop, daemonTrue) accept_thread.start() def stop(self): 停止服务器 self.running False with self.lock: for conn in self.connections.values(): conn.close() self.connections.clear() if self.server_sock: self.server_sock.close() def _accept_loop(self): 接受连接循环 while self.running: try: readable, _, _ select.select([self.server_sock], [], [], 1.0) if readable: client_sock, addr self.server_sock.accept() client_sock.setblocking(False) conn Connection(client_sock, addr) # 启动接收线程 recv_thread threading.Thread( targetself._recv_loop, args(conn,), daemonTrue ) recv_thread.start() # 启动发送线程 send_thread threading.Thread( targetself._send_loop, args(conn,), daemonTrue ) send_thread.start() logger.info(f新连接: {addr}) except Exception as e: if self.running: logger.error(fAccept 错误: {e}) def _recv_loop(self, conn: Connection): 接收循环 buffer b while self.running and not conn.closed: try: readable, _, _ select.select([conn.sock], [], [], 0.1) if readable: data conn.sock.recv(65536) if not data: break buffer data # 解析消息4字节长度前缀 while len(buffer) 4: msg_len struct.unpack_from(!I, buffer, 0)[0] if len(buffer) 4 msg_len: break msg_data buffer[4:4 msg_len] buffer buffer[4 msg_len:] msg Message.deserialize(msg_data) # 记录 node_id 到连接的映射 if msg.sender_id: conn.node_id msg.sender_id with self.lock: self.connections[msg.sender_id] conn # 处理消息 if self.message_handler: self.message_handler(msg, conn) except Exception as e: logger.error(f接收错误: {e}) break self._remove_connection(conn) def _send_loop(self, conn: Connection): 发送循环 while self.running and not conn.closed: try: msg conn.send_queue.get(timeout1.0) conn.send_sync(msg) except queue.Empty: continue except Exception as e: logger.error(f发送循环错误: {e}) break def _remove_connection(self, conn: Connection): 移除连接 with self.lock: for node_id, c in list(self.connections.items()): if c conn: del self.connections[node_id] break conn.close() def send_to(self, node_id: str, msg: Message): 向指定节点发送消息 with self.lock: conn self.connections.get(node_id) if conn and not conn.closed: conn.send(msg) return True return False def broadcast(self, msg: Message, exclude: str None): 广播消息 with self.lock: for node_id, conn in self.connections.items(): if node_id ! exclude and not conn.closed: conn.send(msg) class TCPClient: TCP 客户端用于主动发起连接 def __init__(self): self.connections: Dict[str, Connection] {} self.lock threading.Lock() def connect(self, addr: tuple, node_id: str ) - Optional[Connection]: 连接到远程节点 try: sock socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.settimeout(5.0) sock.connect(addr) sock.setblocking(True) conn Connection(sock, addr, node_id) with self.lock: self.connections[addr] conn logger.info(f连接到 {addr}) return conn except Exception as e: logger.error(f连接 {addr} 失败: {e}) return None def send(self, addr: tuple, msg: Message) - bool: 向指定地址发送消息 with self.lock: conn self.connections.get(addr) if conn and not conn.closed: conn.send_sync(msg) return True return False四、节点发现4.1 基于 Gossip 的节点发现# minikv/discovery/gossip.py import random import time import threading import logging from typing import Set, Dict, List, Optional, Callable from dataclasses import dataclass, field from ..transport.message import Message, MessageType from ..transport.tcp_transport import TCPServer logger logging.getLogger(__name__) dataclass class NodeInfo: 节点信息 node_id: str host: str port: int last_seen: float 0.0 is_alive: bool True metadata: Dict field(default_factorydict) class GossipDiscovery: Gossip 协议节点发现 每个节点定期 1. 随机选择几个邻居 2. 交换各自知道的节点列表 3. 传播新节点和失效节点信息 def __init__(self, node_id: str, host: str, port: int, gossip_interval: float 1.0, fanout: int 3, failure_timeout: float 30.0): self.node_id node_id self.host host self.port port self.gossip_interval gossip_interval self.fanout fanout self.failure_timeout failure_timeout # 已知节点列表 self.nodes: Dict[str, NodeInfo] {} self.local_info NodeInfo( node_idnode_id, hosthost, portport, last_seentime.time() ) self.nodes[node_id] self.local_info # 传输层 self.transport: Optional[TCPServer] None # 回调 self.on_node_join: Optional[Callable] None self.on_node_leave: Optional[Callable] None # 控制 self.running False def start(self, transport: TCPServer): 启动发现服务 self.transport transport self.transport.message_handler self._handle_message self.running True # 启动 Gossip 循环 gossip_thread threading.Thread(targetself._gossip_loop, daemonTrue) gossip_thread.start() # 启动故障检测 failure_thread threading.Thread(targetself._failure_detection_loop, daemonTrue) failure_thread.start() logger.info(fGossip discovery started on {self.host}:{self.port}) def stop(self): 停止发现服务 self.running False def bootstrap(self, seed_nodes: List[tuple]): 加入集群向种子节点发送 JOIN 请求 Args: seed_nodes: [(host, port), ...] for host, port in seed_nodes: addr (host, port) msg Message( msg_typeMessageType.JOIN, sender_idself.node_id, body{ host: self.host, port: self.port, metadata: {} } ) # 通过 TCP 客户端发送 from ..transport.tcp_transport import TCPClient client TCPClient() conn client.connect(addr) if conn: client.send(addr, msg) logger.info(fSent JOIN to {host}:{port}) def _handle_message(self, msg: Message, conn): 处理收到的消息 if msg.msg_type MessageType.PING: self._handle_ping(msg, conn) elif msg.msg_type MessageType.PONG: self._handle_pong(msg, conn) elif msg.msg_type MessageType.JOIN: self._handle_join(msg, conn) elif msg.msg_type MessageType.JOIN_RESPONSE: self._handle_join_response(msg, conn) elif msg.msg_type MessageType.LEAVE: self._handle_leave(msg, conn) def _handle_ping(self, msg: Message, conn): 处理 Ping # 更新节点信息 self._update_node(msg.sender_id, conn) # 回复 Pong带上已知节点列表 known_nodes { nid: { host: info.host, port: info.port, last_seen: info.last_seen, is_alive: info.is_alive } for nid, info in self.nodes.items() if nid ! msg.sender_id } reply msg.reply({known_nodes: known_nodes}) conn.send(reply) def _handle_pong(self, msg: Message, conn): 处理 Pong # 更新节点信息 self._update_node(msg.sender_id, conn) # 合并已知节点 known_nodes msg.body.get(known_nodes, {}) for node_id, info in known_nodes.items(): if node_id not in self.nodes and node_id ! self.node_id: new_node NodeInfo( node_idnode_id, hostinfo[host], portinfo[port], last_seeninfo.get(last_seen, 0), is_aliveinfo.get(is_alive, True) ) self.nodes[node_id] new_node logger.info(f发现新节点: {node_id}{info[host]}:{info[port]}) if self.on_node_join: self.on_node_join(new_node) def _handle_join(self, msg: Message, conn): 处理加入请求 node_info msg.body new_node NodeInfo( node_idmsg.sender_id, hostnode_info[host], portnode_info[port], last_seentime.time() ) self.nodes[msg.sender_id] new_node logger.info(f新节点加入: {msg.sender_id}) # 回复已知节点列表 known_nodes { nid: { host: info.host, port: info.port, last_seen: info.last_seen, is_alive: info.is_alive } for nid, info in self.nodes.items() if nid ! msg.sender_id } reply msg.reply({ status: ok, your_id: self.node_id, known_nodes: known_nodes }) conn.send(reply) if self.on_node_join: self.on_node_join(new_node) def _handle_join_response(self, msg: Message, conn): 处理加入响应 if msg.body.get(status) ok: # 合并已知节点 known_nodes msg.body.get(known_nodes, {}) for node_id, info in known_nodes.items(): if node_id not in self.nodes and node_id ! self.node_id: new_node NodeInfo( node_idnode_id, hostinfo[host], portinfo[port], last_seeninfo.get(last_seen, 0), is_aliveinfo.get(is_alive, True) ) self.nodes[node_id] new_node def _handle_leave(self, msg: Message, conn): 处理离开通知 node_id msg.sender_id if node_id in self.nodes: self.nodes[node_id].is_alive False logger.info(f节点离开: {node_id}) if self.on_node_leave: self.on_node_leave(self.nodes[node_id]) def _gossip_loop(self): Gossip 传播循环 while self.running: time.sleep(self.gossip_interval) # 选择随机邻居 alive_nodes [ nid for nid, info in self.nodes.items() if nid ! self.node_id and info.is_alive ] if not alive_nodes: continue targets random.sample( alive_nodes, min(self.fanout, len(alive_nodes)) ) for target in targets: # 发送 Ping msg Message( msg_typeMessageType.PING, sender_idself.node_id, receiver_idtarget ) self.transport.send_to(target, msg) def _failure_detection_loop(self): 故障检测循环 while self.running: time.sleep(5.0) # 每5秒检查一次 now time.time() for node_id, info in list(self.nodes.items()): if node_id self.node_id: continue if info.is_alive and (now - info.last_seen) self.failure_timeout: info.is_alive False logger.warning(f节点疑似故障: {node_id} (最后活跃: {info.last_seen})) if self.on_node_leave: self.on_node_leave(info) def _update_node(self, node_id: str, conn): 更新节点信息 if node_id not in self.nodes: self.nodes[node_id] NodeInfo( node_idnode_id, hostconn.addr[0], portconn.addr[1] ) self.nodes[node_id].last_seen time.time() self.nodes[node_id].is_alive True self.nodes[node_id].host conn.addr[0] self.nodes[node_id].port conn.addr[1] def get_alive_nodes(self) - List[NodeInfo]: 获取存活节点列表 return [ info for info in self.nodes.values() if info.is_alive and info.node_id ! self.node_id ] def get_node(self, node_id: str) - Optional[NodeInfo]: 获取节点信息 return self.nodes.get(node_id)五、节点实现5.1 完整的节点类# minikv/node.py import logging from typing import List, Optional from .transport.message import Message, MessageType from .transport.tcp_transport import TCPServer from .discovery.gossip import GossipDiscovery, NodeInfo logger logging.getLogger(__name__) class MiniKVNode: MiniKV 节点 def __init__(self, node_id: str, host: str, port: int): self.node_id node_id self.host host self.port port # 传输层 self.transport TCPServer(host, port) # 节点发现 self.discovery GossipDiscovery(node_id, host, port) self.discovery.on_node_join self._on_node_join self.discovery.on_node_leave self._on_node_leave # KV 存储后续实现 self.store {} def start(self): 启动节点 # 启动 TCP 服务器 self.transport.start() # 启动节点发现 self.discovery.start(self.transport) logger.info(fMiniKV 节点启动: {self.node_id}{self.host}:{self.port}) def stop(self): 停止节点 # 发送离开通知 self._broadcast_leave() self.discovery.stop() self.transport.stop() logger.info(fMiniKV 节点停止: {self.node_id}) def join_cluster(self, seed_nodes: List[tuple]): 加入集群 self.discovery.bootstrap(seed_nodes) def _on_node_join(self, node: NodeInfo): 节点加入回调 logger.info(f节点加入集群: {node.node_id}{node.host}:{node.port}) def _on_node_leave(self, node: NodeInfo): 节点离开回调 logger.info(f节点离开集群: {node.node_id}) def _broadcast_leave(self): 广播离开消息 msg Message( msg_typeMessageType.LEAVE, sender_idself.node_id ) self.transport.broadcast(msg) def get_cluster_status(self) - dict: 获取集群状态 alive_nodes self.discovery.get_alive_nodes() return { current_node: self.node_id, alive_nodes: [ { node_id: n.node_id, host: n.host, port: n.port, last_seen: n.last_seen } for n in alive_nodes ], total_nodes: len(alive_nodes) 1 }六、完整演示# examples/cluster_demo.py import time import threading import logging import sys # 配置日志 logging.basicConfig( levellogging.INFO, format%(asctime)s [%(levelname)s] %(name)s: %(message)s ) sys.path.insert(0, ..) from minikv.node import MiniKVNode def start_node(node_id: str, port: int, seed_nodes: list None): 启动一个节点 node MiniKVNode(node_id, localhost, port) node.start() if seed_nodes: time.sleep(0.5) node.join_cluster(seed_nodes) return node def demo(): print( * 60) print( MiniKV 集群启动演示) print( * 60) # 启动三个节点 print(\n 启动节点...) node1 start_node(node-1, 9001) print(f ✅ node-1 启动 (port 9001)) time.sleep(0.5) node2 start_node(node-2, 9002, [(localhost, 9001)]) print(f ✅ node-2 启动 (port 9002, 加入 node-1)) time.sleep(0.5) node3 start_node(node-3, 9003, [(localhost, 9001)]) print(f ✅ node-3 启动 (port 9003, 加入 node-1)) # 等待 Gossip 传播 print(\n⏳ 等待 Gossip 传播...) time.sleep(3) # 查看集群状态 print(\n 集群状态:) for i, node in enumerate([node1, node2, node3], 1): status node.get_cluster_status() print(f\n Node{i} ({node.node_id}) 视角:) print(f 存活节点: {status[total_nodes]}) for n in status[alive_nodes]: print(f - {n[node_id]} {n[host]}:{n[port]}) # 测试节点故障 print(\n 模拟节点故障: 停止 node-3...) node3.stop() time.sleep(2) print(\n 故障后集群状态 (node-1 视角):) status node1.get_cluster_status() print(f 存活节点: {status[total_nodes]}) for n in status[alive_nodes]: alive ✅ if n[node_id] ! node-3 else ❌ print(f {alive} {n[node_id]} {n[host]}:{n[port]}) # 清理 print(\n 清理...) node1.stop() node2.stop() print(\n✅ 演示完成) if __name__ __main__: demo()七、测试# tests/test_transport.py import unittest import threading import time from minikv.transport.message import Message, MessageType from minikv.transport.tcp_transport import TCPServer, TCPClient class TestTransport(unittest.TestCase): 传输层测试 def test_message_serialization(self): 测试消息序列化 msg Message( msg_typeMessageType.PING, sender_idnode-1, receiver_idnode-2, body{key: value} ) data msg.serialize() restored Message.deserialize(data) self.assertEqual(msg.msg_id, restored.msg_id) self.assertEqual(msg.msg_type, restored.msg_type) self.assertEqual(msg.sender_id, restored.sender_id) self.assertEqual(msg.body, restored.body) def test_reply_generation(self): 测试回复生成 request Message( msg_typeMessageType.PING, sender_idnode-1, receiver_idnode-2 ) reply request.reply({status: ok}) self.assertEqual(reply.msg_type, MessageType.PONG) self.assertEqual(reply.sender_id, node-2) self.assertEqual(reply.receiver_id, node-1) self.assertEqual(reply.body[reply_to], request.msg_id) def test_tcp_communication(self): 测试 TCP 通信 received [] def handler(msg, conn): received.append(msg) # 启动服务器 server TCPServer(localhost, 9999) server.message_handler handler server.start() time.sleep(0.2) # 客户端连接 client TCPClient() conn client.connect((localhost, 9999)) # 发送消息 msg Message( msg_typeMessageType.PING, sender_idtest-client, body{hello: world} ) client.send((localhost, 9999), msg) time.sleep(0.2) self.assertEqual(len(received), 1) self.assertEqual(received[0].msg_type, MessageType.PING) self.assertEqual(received[0].body[hello], world) server.stop() class TestGossipDiscovery(unittest.TestCase): Gossip 发现测试 def test_node_discovery(self): 测试节点发现 nodes [] for i in range(3): from minikv.node import MiniKVNode node MiniKVNode(fnode-{i}, localhost, 9010 i) node.start() nodes.append(node) time.sleep(0.5) # node-1 和 node-2 加入 node-0 nodes[1].join_cluster([(localhost, 9010)]) nodes[2].join_cluster([(localhost, 9010)]) # 等待 Gossip 传播 time.sleep(3) # 验证所有节点都知道彼此 for node in nodes: status node.get_cluster_status() self.assertEqual(status[total_nodes], 3) # 清理 for node in nodes: node.stop() if __name__ __main__: unittest.main()八、总结这一讲我们实现了 MiniKV 的分布式基石组件功能消息系统​消息类型定义、JSON序列化、二进制协议TCP传输​非阻塞IO、连接管理、收发队列Gossip发现​节点自发现、故障检测、信息传播节点实现​完整的节点生命周期管理现在三个节点可以自动发现彼此、感知故障、传播集群状态。这为后续实现 Raft 共识和数据分片打下了基础。下一讲我们将实现一致性哈希让数据均匀分布到集群节点上并支持动态扩缩容。 开发之余处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top子页 PDF 大师PDF 大师 - zz365工具箱。所有计算在浏览器完成文件不上传服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。
返回列表