【问题标题】:What's the best pattern for a thread safe write cache to database?线程安全写入缓存到数据库的最佳模式是什么?
【发布时间】:2020-11-13 08:28:12
【问题描述】:

我有一个可以被多个线程调用的方法,用于将数据写入数据库。为了减少数据库流量,我缓存数据并批量写入。

现在我想知道,有没有更好的(例如无锁模式)可以使用?

这是我目前如何做的示例?

    public class WriteToDatabase : IWriter, IDisposable
    {
        public WriteToDatabase(PLCProtocolServiceConfig currentConfig)
        {
            writeTimer = new System.Threading.Timer(Writer);
            writeTimer.Change((int)currentConfig.WriteToDatabaseTimer.TotalMilliseconds, Timeout.Infinite);
            this.currentConfig = currentConfig;
        }

        private System.Threading.Timer writeTimer;
        private List<PlcProtocolDTO> writeChache = new List<PlcProtocolDTO>();
        private readonly PLCProtocolServiceConfig currentConfig;
        private bool disposed;

        public void Write(PlcProtocolDTO row)
        {
            lock (this)
            {
                writeChache.Add(row);
            }
        }

        private void Writer(object state)
        {
            List<PlcProtocolDTO> oldCachce = null;
            lock (this)
            {
                if (writeChache.Count > 0)
                {
                    oldCachce = writeChache;
                    writeChache = new List<PlcProtocolDTO>();
                }
            }

            if (oldCachce != null)
            {
                    using (var s = VisuDL.CreateSession())
                    {
                        s.Insert(oldCachce);
                    }
            }

            if (!this.disposed)
                writeTimer.Change((int)currentConfig.WriteToDatabaseTimer.TotalMilliseconds, Timeout.Infinite);
        }

        public void Dispose()
        {
            this.disposed = true;
            writeTimer.Dispose();
            Writer(null);
        }
    }

【问题讨论】:

标签: c# multithreading caching design-patterns


【解决方案1】:

我可以看到基于计时器的代码存在一些问题。

  • 即使在新版本的代码中,仍然有可能在重新启动或关闭时丢失写入。 Dispose 方法不等待当前可能正在进行的最后一个计时器回调完成。 由于计时器回调在线程池线程(后台线程)上运行,因此它们将在主线程退出时中止。
  • 批次的大小没有限制,当您达到底层存储 API 的限制时,这将被打破 (例如 sql 数据库对查询长度和使用的参数数量有限制)。
  • 由于您正在执行 i/o,因此实现可能应该是异步的
  • 这将在负载下表现不佳。 特别是随着负载不断增加,批次会变大,因此执行速度会变慢, 反过来,较慢的批处理执行将给下一个额外的时间来积累项目,使它们变得更慢,等等...... 最终要么写入批处理将失败(如果您达到 sql 限制或查询超时),或者应用程序将内存不足。 要处理高负载,您实际上只有两个选择,即应用背压(即减慢生产者的速度)或放弃写入。
  • 如果数据库可以处理,您可能希望允许有限数量的并发写入者。
  • disposed 字段存在竞争条件,这可能会导致writeTimer.Change 中的ObjectDisposedException

我认为解决上述问题的更好模式是消费者-生产者模式,您可以在 .net 中实现它 使用 ConcurrentQueue 或使用新的 System.Threading.Channels api。

另外请记住,如果您的应用程序因任何原因崩溃,您将丢失仍在缓冲中的记录。

这是一个使用通道的示例实现:

public interface IWriter<in T>
{
    ValueTask WriteAsync(IEnumerable<T> items);
}

public sealed record Options(int BatchSize, TimeSpan Interval, int MaxPendingWrites, int Concurrency);

public class BatchWriter<T> : IWriter<T>, IAsyncDisposable
{
    readonly IWriter<T> writer;
    readonly Options options;
    readonly Channel<T> channel;
    readonly Task[] consumers;

    public BatchWriter(IWriter<T> writer, Options options)
    {
        this.writer = writer;
        this.options = options;

        channel = Channel.CreateBounded<T>(new BoundedChannelOptions(options.MaxPendingWrites)
        {
            // Choose between backpressure (Wait) or
            // various ways to drop writes (DropNewest, DropOldest, DropWrite).
            FullMode = BoundedChannelFullMode.Wait,

            SingleWriter = false,
            SingleReader = options.Concurrency == 1
        });

        consumers = Enumerable.Range(start: 0, options.Concurrency)
            .Select(_ => Task.Run(Start))
            .ToArray();
    }

    async Task Start()
    {
        var batch = new List<T>(options.BatchSize);

        var timer = Task.Delay(options.Interval);
        var canRead = channel.Reader.WaitToReadAsync().AsTask();

        while (true)
        {
            if (await Task.WhenAny(timer, canRead) == timer)
            {
                timer = Task.Delay(options.Interval);
                await Flush(batch);
            }
            else if (await canRead)
            {
                while (channel.Reader.TryRead(out var item))
                {
                    batch.Add(item);

                    if (batch.Count == options.BatchSize)
                    {
                        await Flush(batch);
                    }
                }

                canRead = channel.Reader.WaitToReadAsync().AsTask();
            }
            else
            {
                await Flush(batch);
                return;
            }
        }

        async Task Flush(ICollection<T> items)
        {
            if (items.Count > 0)
            {
                await writer.WriteAsync(items);
                items.Clear();
            }
        }
    }

    public async ValueTask WriteAsync(IEnumerable<T> items)
    {
        foreach (var item in items)
        {
            await channel.Writer.WriteAsync(item);
        }
    }

    public async ValueTask DisposeAsync()
    {
        channel.Writer.Complete();
        await Task.WhenAll(consumers);
    }
}

【讨论】:

  • 这只是一个简化的示例,以获取想法 :-) 我们的代码中有多个位置,我们收集一些更改并批量处理它们。实际上,当 Timer Callback 完成时,我再次启动 Timer。我还通过 dispose 停止计时器并将缓存写入数据库。
  • @user1237393 你是否也在限制批量大小和限制生产者?
【解决方案2】:

您可以使用ImmutableList,而不是使用可变的List 并使用锁来保护它,并且不必担心列表在错误的时间被错误的线程改变的可能性。使用immutable collections,传递数据快照既便宜又容易,因为您在创建数据副本时不需要阻止写入者(可能还有读取者)。不可变集合本身就是一个快照。

虽然您不必担心集合的内容,但您仍然需要担心它的引用。这是因为更新不可变集合意味着用新集合替换对旧集合的引用。您不希望多个线程以无法控制的方式交换引用,因此您仍然需要某种同步。您仍然可以使用locks,但是通过使用互锁操作很容易完全避免锁定。下面的示例使用方便的ImmutableInterlocked.Update 方法,它允许在一行中进行原子更新和交换:

private ImmutableList<PlcProtocolDTO> writeCache
    = ImmutableList<PlcProtocolDTO>.Empty;

public void Write(PlcProtocolDTO row)
{
    ImmutableInterlocked.Update(ref writeCache, x => x.Add(row));
}

private void Writer(object state)
{
    IList<PlcProtocolDTO> oldCache = Interlocked.Exchange(
        ref writeCache, ImmutableList<PlcProtocolDTO>.Empty);

    using (var s = VisuDL.CreateSession())
        s.Insert(oldCache);
}

private void Dump()
{
    foreach (var row in Volatile.Read(ref writeCache))
        Console.WriteLine(row);
}

这里是ImmutableInterlocked.Update方法的描述:

通过指定的转换函数使用乐观锁定事务语义就地改变值。转换会根据需要重试多次,以赢得乐观锁定竞赛。

此方法可用于更新任何类型的引用类型变量。随着 new C# 9 record types 的出现,它的使用量可能会增加,默认情况下它是不可变的,并且打算按原样使用。

【讨论】:

  • 您认为 ImmutableList 是否比 Lock 更具性能?我不知道,但是有了这个,每次都会重新创建列表......也许我需要用基准 dotnet 做一些测试
  • 没有写入丢失是否安全?为什么我还需要 ImmutableList?
  • @user1237393 经常创建新的不可变列表是负担得起的,因为它们在内部实现为二叉树,主要由可重用的节点组成。与普通或并发集合上的相同操作相比,不可变集合上的基本操作的性能不是很好(它们通常至少慢 10 倍)。每当您需要拍摄数据快照时,回报就会到来。如果您从不需要快照,那么就没有性能回报,剩下的唯一优势就是架构优势(当然这是主观的)。
  • @user1237393 如果您正确使用互锁操作或锁定,则不会丢失任何更新。对于(悲观的)lock,您需要保护添加到集合和引用交换。使用(乐观的)ImmutableInterlocked.Update,您可以确保在与另一个线程竞争的情况下,更新操作将在赢得比赛的线程先前返回的版本上重复。省略锁的代价是可能会旋转一点(并为 GC 制造一些垃圾)。无论如何,更新不可能逃脱。
猜你喜欢
  • 2016-04-17
  • 2021-05-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多