You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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静态方法,避免重复传入委托:

  1. 定义标记接口并修改基类:
public interface IStreamEntryFactory { }

public abstract class RedisStreamEntryDto : IStreamEntryFactory
{
    public required float Progress { get; set; }
    public required bool RunDone { get; set; }
}
  1. 修改泛型方法:
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:抽象实例方法+无参构造

在基类定义抽象加载方法,先实例化子类再加载数据,要求子类有无参构造函数:

  1. 修改基类:
public abstract class RedisStreamEntryDto
{
    public required float Progress { get; set; }
    public required bool RunDone { get; set; }

    // 抽象加载方法,子类实现具体逻辑
    public abstract void LoadFromStreamEntry(StreamEntry entry);
}
  1. 修改泛型方法:
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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 07:44:57