dotNext Raft多进程下持久化集群配置失败及Announcer使用咨询
dotNext Raft 动态节点加入与共享配置存储解决方案
问题描述
我正在研究dotNext中的Raft组件,想把启动时预注册所有集群节点的方式,改成用ClusterMemberAnnouncer实现新节点加入时自动通知领导者。按照理解,初始节点用ColdStart模式启动,后续节点通过Announcer加入集群,已经写了Announcer的注册代码:
services.AddTransient<ClusterMemberAnnouncer<UriEndPoint>>(serviceProvider => async (memberId, address, cancellationToken) => { // Register the node with the configuration storage var configurationStorage = serviceProvider.GetService<IClusterConfigurationStorage<UriEndPoint>>(); if (configurationStorage == null) throw new Exception("Unable to resolve the IClusterConfigurationStorage when adding the new node member"); await configurationStorage.AddMemberAsync(memberId, address, cancellationToken); });
但使用文档中的services.UsePersistentConfigurationStorage("configurationStorage")本地文件存储时,多进程(独立控制台窗口)启动节点会报错:
The process cannot access the file 'C:\Projects\RaftTest\configurationStorage\active.list' because it is being used by another process.
需要两个帮助:
- dotNext Raft中使用Announcer的完整示例
- 多进程/Docker容器环境下可用的共享持久化集群配置存储实现(附示例)
一、完整的Announcer使用示例
1. 初始节点配置(ColdStart模式)
初始节点需要以冷启动模式初始化集群,同时配置Announcer和共享存储:
var services = new ServiceCollection(); // 配置共享集群配置存储(下文会实现Redis版本) services.AddSingleton<IClusterConfigurationStorage<UriEndPoint>, RedisClusterConfigurationStorage>(); // 注册Announcer services.AddTransient<ClusterMemberAnnouncer<UriEndPoint>>(sp => async (memberId, address, ct) => { var storage = sp.GetRequiredService<IClusterConfigurationStorage<UriEndPoint>>(); await storage.AddMemberAsync(memberId, address, ct); }); // 配置Raft节点 services.AddRaftCluster<UriEndPoint>(builder => { builder.ConfigureCluster(options => { options.ClusterId = "my-raft-cluster"; options.MemberId = Guid.NewGuid(); // 初始节点的唯一ID options.EndPoint = new UriEndPoint(new Uri("http://localhost:5000")); }) .UseColdStart() // 冷启动创建新集群 .UseHttpTransport() .UseMetadataStorage(); }); // 构建服务并启动 var serviceProvider = services.BuildServiceProvider(); var clusterHost = serviceProvider.GetRequiredService<IHostedService>(); await clusterHost.StartAsync(CancellationToken.None);
2. 后续节点配置(加入现有集群)
后续节点不需要ColdStart,只需配置自身地址并依赖Announcer完成注册:
var services = new ServiceCollection(); // 同样配置共享存储 services.AddSingleton<IClusterConfigurationStorage<UriEndPoint>, RedisClusterConfigurationStorage>(); // 注册相同的Announcer services.AddTransient<ClusterMemberAnnouncer<UriEndPoint>>(sp => async (memberId, address, ct) => { var storage = sp.GetRequiredService<IClusterConfigurationStorage<UriEndPoint>>(); await storage.AddMemberAsync(memberId, address, ct); }); // 配置Raft节点(无ColdStart) services.AddRaftCluster<UriEndPoint>(builder => { builder.ConfigureCluster(options => { options.ClusterId = "my-raft-cluster"; // 必须和初始集群一致 options.MemberId = Guid.NewGuid(); // 新节点的唯一ID options.EndPoint = new UriEndPoint(new Uri("http://localhost:5001")); }) .UseHttpTransport() .UseMetadataStorage(); }); // 启动节点 var serviceProvider = services.BuildServiceProvider(); var clusterHost = serviceProvider.GetRequiredService<IHostedService>(); await clusterHost.StartAsync(CancellationToken.None);
二、共享持久化配置存储实现(Redis示例)
默认的PersistentConfigurationStorage是本地文件系统实现,不支持多进程/容器共享。下面是基于Redis的IClusterConfigurationStorage实现,适合分布式环境:
using DotNext.Net.Cluster.Consensus.Raft; using StackExchange.Redis; public class RedisClusterConfigurationStorage : IClusterConfigurationStorage<UriEndPoint> { private readonly IDatabase _redisDb; private const string ActiveMembersKey = "raft:active-members"; private const string PassiveMembersKey = "raft:passive-members"; public RedisClusterConfigurationStorage(IConnectionMultiplexer redisMultiplexer) { _redisDb = redisMultiplexer.GetDatabase(); } public async ValueTask<IReadOnlyDictionary<ClusterMemberId, UriEndPoint>> GetActiveMembersAsync(CancellationToken token) { var entries = await _redisDb.HashGetAllAsync(ActiveMembersKey); return entries.ToDictionary( entry => new ClusterMemberId(Guid.Parse(entry.Name)), entry => new UriEndPoint(new Uri(entry.Value)) ); } public async ValueTask<IReadOnlyDictionary<ClusterMemberId, UriEndPoint>> GetPassiveMembersAsync(CancellationToken token) { var entries = await _redisDb.HashGetAllAsync(PassiveMembersKey); return entries.ToDictionary( entry => new ClusterMemberId(Guid.Parse(entry.Name)), entry => new UriEndPoint(new Uri(entry.Value)) ); } public async ValueTask AddMemberAsync(ClusterMemberId memberId, UriEndPoint address, CancellationToken token) { await _redisDb.HashSetAsync(ActiveMembersKey, memberId.ToString(), address.ToString()); } public async ValueTask RemoveMemberAsync(ClusterMemberId memberId, CancellationToken token) { await _redisDb.HashDeleteAsync(ActiveMembersKey, memberId.ToString()); await _redisDb.HashDeleteAsync(PassiveMembersKey, memberId.ToString()); } public async ValueTask PromoteToActiveAsync(ClusterMemberId memberId, CancellationToken token) { var value = await _redisDb.HashGetAsync(PassiveMembersKey, memberId.ToString()); if (!value.IsNullOrEmpty) { await _redisDb.HashSetAsync(ActiveMembersKey, memberId.ToString(), value); await _redisDb.HashDeleteAsync(PassiveMembersKey, memberId.ToString()); } } public async ValueTask DemoteToPassiveAsync(ClusterMemberId memberId, UriEndPoint address, CancellationToken token) { await _redisDb.HashSetAsync(PassiveMembersKey, memberId.ToString(), address.ToString()); await _redisDb.HashDeleteAsync(ActiveMembersKey, memberId.ToString()); } }
注册Redis依赖
在服务配置中添加Redis连接:
// 配置Redis连接 services.AddSingleton<IConnectionMultiplexer>(sp => { var config = ConfigurationOptions.Parse("localhost:6379"); return ConnectionMultiplexer.Connect(config); }); // 注册自定义存储 services.AddSingleton<IClusterConfigurationStorage<UriEndPoint>, RedisClusterConfigurationStorage>();
内容的提问来源于stack exchange,提问作者James_2195
相关产品推荐
相关产品推荐

