【问题标题】:Bandwidth throttling in C#C#中的带宽限制
【发布时间】:2010-09-27 03:04:31
【问题描述】:

我正在开发一个在后台持续发送数据流的程序,我希望允许用户设置上传和下载限制的上限。

我已经阅读了 token bucket 和 leaky bucket 算法,似乎后者似乎符合描述,因为这不是最大化网络带宽的问题,而是尽可能不引人注目。

然而,我有点不确定我将如何实现这一点。一种自然的方法是扩展抽象 Stream 类以简化现有流量的扩展,但这是否不需要额外的线程参与来发送数据同时接收(漏桶)?任何有关其他实现相同功能的提示将不胜感激。

另外,虽然我可以修改程序接收的数据量,但带宽限制在 C# 级别的工作情况如何?计算机是否仍会接收数据并简单地保存数据,从而有效地取消节流效果,还是会等到我要求接收更多数据?

编辑:我对限制传入和传出数据感兴趣,我无法控制流的另一端。

【问题讨论】:

    标签: c# bandwidth throttling


    【解决方案1】:

    我想出了 arul 提到的 ThrottledStream-Class 的不同实现。我的版本使用 WaitHandle 和间隔为 1 秒的计时器:

    public ThrottledStream(Stream parentStream, int maxBytesPerSecond=int.MaxValue) 
    {
        MaxBytesPerSecond = maxBytesPerSecond;
        parent = parentStream;
        processed = 0;
        resettimer = new System.Timers.Timer();
        resettimer.Interval = 1000;
        resettimer.Elapsed += resettimer_Elapsed;
        resettimer.Start();         
    }
    
    protected void Throttle(int bytes)
    {
        try
        {
            processed += bytes;
            if (processed >= maxBytesPerSecond)
                wh.WaitOne();
        }
        catch
        {
        }
    }
    
    private void resettimer_Elapsed(object sender, ElapsedEventArgs e)
    {
        processed = 0;
        wh.Set();
    }
    

    只要带宽限制超过线程将休眠直到下一秒开始。无需计算最佳睡眠时长。

    全面实施:

    public class ThrottledStream : Stream
    {
        #region Properties
    
        private int maxBytesPerSecond;
        /// <summary>
        /// Number of Bytes that are allowed per second
        /// </summary>
        public int MaxBytesPerSecond
        {
            get { return maxBytesPerSecond; }
            set 
            {
                if (value < 1)
                    throw new ArgumentException("MaxBytesPerSecond has to be >0");
    
                maxBytesPerSecond = value; 
            }
        }
    
        #endregion
    
    
        #region Private Members
    
        private int processed;
        System.Timers.Timer resettimer;
        AutoResetEvent wh = new AutoResetEvent(true);
        private Stream parent;
    
        #endregion
    
        /// <summary>
        /// Creates a new Stream with Databandwith cap
        /// </summary>
        /// <param name="parentStream"></param>
        /// <param name="maxBytesPerSecond"></param>
        public ThrottledStream(Stream parentStream, int maxBytesPerSecond=int.MaxValue) 
        {
            MaxBytesPerSecond = maxBytesPerSecond;
            parent = parentStream;
            processed = 0;
            resettimer = new System.Timers.Timer();
            resettimer.Interval = 1000;
            resettimer.Elapsed += resettimer_Elapsed;
            resettimer.Start();         
        }
    
        protected void Throttle(int bytes)
        {
            try
            {
                processed += bytes;
                if (processed >= maxBytesPerSecond)
                    wh.WaitOne();
            }
            catch
            {
            }
        }
    
        private void resettimer_Elapsed(object sender, ElapsedEventArgs e)
        {
            processed = 0;
            wh.Set();
        }
    
        #region Stream-Overrides
    
        public override void Close()
        {
            resettimer.Stop();
            resettimer.Close();
            base.Close();
        }
        protected override void Dispose(bool disposing)
        {
            resettimer.Dispose();
            base.Dispose(disposing);
        }
    
        public override bool CanRead
        {
            get { return parent.CanRead; }
        }
    
        public override bool CanSeek
        {
            get { return parent.CanSeek; }
        }
    
        public override bool CanWrite
        {
            get { return parent.CanWrite; }
        }
    
        public override void Flush()
        {
            parent.Flush();
        }
    
        public override long Length
        {
            get { return parent.Length; }
        }
    
        public override long Position
        {
            get
            {
                return parent.Position;
            }
            set
            {
                parent.Position = value;
            }
        }
    
        public override int Read(byte[] buffer, int offset, int count)
        {
            Throttle(count);
            return parent.Read(buffer, offset, count);
        }
    
        public override long Seek(long offset, SeekOrigin origin)
        {
            return parent.Seek(offset, origin);
        }
    
        public override void SetLength(long value)
        {
            parent.SetLength(value);
        }
    
        public override void Write(byte[] buffer, int offset, int count)
        {
            Throttle(count);
            parent.Write(buffer, offset, count);
        }
    
        #endregion
    
    
    }
    

    【讨论】:

    • 如果在计时器计时时不将processed设置为0,而是从中减去maxBytesPerSecond,它会变得更准确。
    • 在读取中,这会比限制慢。例如,您从 Internet 下载。缓冲区 8Kib,每次读取速度 1Kib,限制 1Mib/sec。然后你每次读取损失 7Kib,它在第 128 次读取时wh.WaitOne() -> 实际速度为 16Kib/秒。需要修复int read = parent.Read(buffer, offset, count); Throttle(read); return read;
    【解决方案2】:

    基于@0xDEADBEEF 的解决方案,我创建了以下基于 Rx 调度程序的(可测试的)解决方案:

    public class ThrottledStream : Stream
    {
        private readonly Stream parent;
        private readonly int maxBytesPerSecond;
        private readonly IScheduler scheduler;
        private readonly IStopwatch stopwatch;
    
        private long processed;
    
        public ThrottledStream(Stream parent, int maxBytesPerSecond, IScheduler scheduler)
        {
            this.maxBytesPerSecond = maxBytesPerSecond;
            this.parent = parent;
            this.scheduler = scheduler;
            stopwatch = scheduler.StartStopwatch();
            processed = 0;
        }
    
        public ThrottledStream(Stream parent, int maxBytesPerSecond)
            : this (parent, maxBytesPerSecond, Scheduler.Immediate)
        {
        }
    
        protected void Throttle(int bytes)
        {
            processed += bytes;
            var targetTime = TimeSpan.FromSeconds((double)processed / maxBytesPerSecond);
            var actualTime = stopwatch.Elapsed;
            var sleep = targetTime - actualTime;
            if (sleep > TimeSpan.Zero)
            {
                using (var waitHandle = new AutoResetEvent(initialState: false))
                {
                    scheduler.Sleep(sleep).GetAwaiter().OnCompleted(() => waitHandle.Set());
                    waitHandle.WaitOne();
                }
            }
        }
    
        public override bool CanRead
        {
            get { return parent.CanRead; }
        }
    
        public override bool CanSeek
        {
            get { return parent.CanSeek; }
        }
    
        public override bool CanWrite
        {
            get { return parent.CanWrite; }
        }
    
        public override void Flush()
        {
            parent.Flush();
        }
    
        public override long Length
        {
            get { return parent.Length; }
        }
    
        public override long Position
        {
            get
            {
                return parent.Position;
            }
            set
            {
                parent.Position = value;
            }
        }
    
        public override int Read(byte[] buffer, int offset, int count)
        {
            var read = parent.Read(buffer, offset, count);
            Throttle(read);
            return read;
        }
    
        public override long Seek(long offset, SeekOrigin origin)
        {
            return parent.Seek(offset, origin);
        }
    
        public override void SetLength(long value)
        {
            parent.SetLength(value);
        }
    
        public override void Write(byte[] buffer, int offset, int count)
        {
            Throttle(count);
            parent.Write(buffer, offset, count);
        }
    }
    

    还有一些只需要几毫秒的测试:

    [TestMethod]
    public void ShouldThrottleReading()
    {
        var content = Enumerable
            .Range(0, 1024 * 1024)
            .Select(_ => (byte)'a')
            .ToArray();
        var scheduler = new TestScheduler();
        var source = new ThrottledStream(new MemoryStream(content), content.Length / 8, scheduler);
        var target = new MemoryStream();
    
        var t = source.CopyToAsync(target);
    
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(4).Ticks);
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(8).Ticks - 1);
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(8).Ticks);
        t.Wait(10).Should().BeTrue();
    }
    
    [TestMethod]
    public void ShouldThrottleWriting()
    {
        var content = Enumerable
            .Range(0, 1024 * 1024)
            .Select(_ => (byte)'a')
            .ToArray();
        var scheduler = new TestScheduler();
        var source = new MemoryStream(content);
        var target = new ThrottledStream(new MemoryStream(), content.Length / 8, scheduler);
    
        var t = source.CopyToAsync(target);
    
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(4).Ticks);
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(8).Ticks - 1);
        t.Wait(10).Should().BeFalse();
        scheduler.AdvanceTo(TimeSpan.FromSeconds(8).Ticks);
        t.Wait(10).Should().BeTrue();
    }
    

    【讨论】:

    • 工作,谢谢!那些想要使用它的人,需要包括 Nuget 的 Reactive Extensions
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-08-09
    相关资源
    最近更新 更多