C#中如何实现基于抽象类的Redis流条目对象实例化?
Redis流条目实例化问题解决方案
问题描述
我需要实现一个监听Redis流的方法,根据流类型返回对应条目列表。已定义抽象基类RedisStreamEntryDto,包含所有流条目的公共属性:
public abstract class RedisStreamEntryDto { public required float Progress { get; set; } public required bool RunDone { get; set; } }
监听流的RedisUtils类代码如下,核心问题是如何从StreamEntry实例化对应的RedisStreamEntryDto子类:
public class RedisUtils(ConnectionMultiplexer redis) { public async Task<List<RedisStreamEntryDto>> GetStreamEntriesAsync(string streamName) { var db = redis.GetDatabase(); bool stopListening = false; List<RedisStreamEntryDto> streamItems = []; while (!stopListening) { StreamEntry[] streamEntries = await db.StreamReadAsync(streamName, "0-0"); foreach (var entry in streamEntries) { if(entry.IsNull) { continue; } var add_entry = // 如何实例化对应的子类? streamItems.Add(add_entry); if (add_entry.RunDone) { stopListening = true; break; } } } return streamItems; } }
通常我会用JSON反序列化处理,但收到的是RedisValue类型,没有内置反序列化能力,转成JSON字符串又不够优雅。我曾想让所有子类实现静态方法FromStreamEntry(StreamEntry entry),但抽象类和接口都无法定义抽象静态方法。目前想到的只有单独的工厂类,但觉得太繁琐,有没有更优的实现方式?
另外,我知道调用方期望的条目类型,所以可以改成泛型版本,但问题依然存在:
public class RedisUtils(ConnectionMultiplexer redis) { public async Task<List<T>> GetStreamEntriesAsync<T>(string streamName) where T: RedisStreamEntryDto { var db = redis.GetDatabase(); bool stopListening = false; List<T> streamItems = []; while (!stopListening) { StreamEntry[] streamEntries = await db.StreamReadAsync(streamName, "0-0"); foreach (var entry in streamEntries) { if(entry.IsNull) { continue; } var add_entry = T.FromStreamEntry(entry); // 无法直接调用静态方法 streamItems.Add(add_entry); if (add_entry.RunDone) { stopListening = true; break; } } } return streamItems; } }
可行方案
方案1:传入实例化委托(最灵活)
直接在泛型方法中传入Func<StreamEntry, T>委托,让调用方负责提供实例化逻辑,无需修改基类或子类:
public class RedisUtils(ConnectionMultiplexer redis) { public async Task<List<T>> GetStreamEntriesAsync<T>(string streamName, Func<StreamEntry, T> entryFactory) where T : RedisStreamEntryDto { var db = redis.GetDatabase(); bool stopListening = false; List<T> streamItems = []; while (!stopListening) { StreamEntry[] streamEntries = await db.StreamReadAsync(streamName, "0-0"); foreach (var entry in streamEntries) { if(entry.IsNull) { continue; } var add_entry = entryFactory(entry); streamItems.Add(add_entry); if (add_entry.RunDone) { stopListening = true; break; } } } return streamItems; } }
调用示例:
// 假设子类是TaskStreamEntryDto,已实现FromStreamEntry静态方法 var entries = await redisUtils.GetStreamEntriesAsync("task-stream", TaskStreamEntryDto.FromStreamEntry);
方案2:接口标记+反射(自动绑定静态方法)
定义标记接口,通过反射自动调用子类的FromStreamEntry静态方法,避免重复传入委托:
- 定义标记接口并修改基类:
public interface IStreamEntryFactory { } public abstract class RedisStreamEntryDto : IStreamEntryFactory { public required float Progress { get; set; } public required bool RunDone { get; set; } }
- 修改泛型方法:
public class RedisUtils(ConnectionMultiplexer redis) { public async Task<List<T>> GetStreamEntriesAsync<T>(string streamName) where T : RedisStreamEntryDto, IStreamEntryFactory { var db = redis.GetDatabase(); bool stopListening = false; List<T> streamItems = []; // 提前获取静态方法,避免循环内重复反射 var factoryMethod = typeof(T).GetMethod("FromStreamEntry", new[] { typeof(StreamEntry) }); if (factoryMethod == null || !factoryMethod.IsStatic) { throw new InvalidOperationException($"{typeof(T)}必须实现静态方法FromStreamEntry(StreamEntry)"); } while (!stopListening) { StreamEntry[] streamEntries = await db.StreamReadAsync(streamName, "0-0"); foreach (var entry in streamEntries) { if(entry.IsNull) { continue; } var add_entry = (T)factoryMethod.Invoke(null, new object[] { entry }); streamItems.Add(add_entry); if (add_entry.RunDone) { stopListening = true; break; } } } return streamItems; } }
子类实现示例:
public class TaskStreamEntryDto : RedisStreamEntryDto { public string TaskId { get; set; } public static TaskStreamEntryDto FromStreamEntry(StreamEntry entry) { return new TaskStreamEntryDto { Progress = float.Parse(entry.Values.First(v => v.Name == "Progress").Value), RunDone = bool.Parse(entry.Values.First(v => v.Name == "RunDone").Value), TaskId = entry.Values.First(v => v.Name == "TaskId").Value }; } }
方案3:抽象实例方法+无参构造
在基类定义抽象加载方法,先实例化子类再加载数据,要求子类有无参构造函数:
- 修改基类:
public abstract class RedisStreamEntryDto { public required float Progress { get; set; } public required bool RunDone { get; set; } // 抽象加载方法,子类实现具体逻辑 public abstract void LoadFromStreamEntry(StreamEntry entry); }
- 修改泛型方法:
public class RedisUtils(ConnectionMultiplexer redis) { public async Task<List<T>> GetStreamEntriesAsync<T>(string streamName) where T : RedisStreamEntryDto, new() { var db = redis.GetDatabase(); bool stopListening = false; List<T> streamItems = []; while (!stopListening) { StreamEntry[] streamEntries = await db.StreamReadAsync(streamName, "0-0"); foreach (var entry in streamEntries) { if(entry.IsNull) { continue; } var add_entry = new T(); add_entry.LoadFromStreamEntry(entry); streamItems.Add(add_entry); if (add_entry.RunDone) { stopListening = true; break; } } } return streamItems; } }
子类实现示例:
public class TaskStreamEntryDto : RedisStreamEntryDto { public string TaskId { get; set; } public override void LoadFromStreamEntry(StreamEntry entry) { Progress = float.Parse(entry.Values.First(v => v.Name == "Progress").Value); RunDone = bool.Parse(entry.Values.First(v => v.Name == "RunDone").Value); TaskId = entry.Values.First(v => v.Name == "TaskId").Value; } }
方案对比
- 方案1:性能最优、灵活性最高,无需修改现有类结构,推荐作为首选。
- 方案2:无需传入委托,但依赖反射,性能略低,适合调用方较多的场景。
- 方案3:结构清晰,但要求子类支持无参构造,实例化后再加载数据的模式部分场景下不够直观。
内容的提问来源于stack exchange,提问作者Roland Deschain
相关产品推荐
相关产品推荐

