diff --git a/src/Orleans/Providers/IOrleansProvider.cs b/src/Orleans/Providers/IOrleansProvider.cs index 575d1fc2aed..f86f8b9a778 100644 --- a/src/Orleans/Providers/IOrleansProvider.cs +++ b/src/Orleans/Providers/IOrleansProvider.cs @@ -104,6 +104,17 @@ public static int GetIntProperty(this IProviderConfiguration config, string key, return config.Properties.TryGetValue(key, out s) ? int.Parse(s) : settingDefault; } + public static bool TryGetDoubleProperty(this IProviderConfiguration config, string key, out double setting) + { + if (config == null) + { + throw new ArgumentNullException("config"); + } + string s; + setting = 0; + return config.Properties.TryGetValue(key, out s) ? double.TryParse(s, out setting) : false; + } + public static string GetProperty(this IProviderConfiguration config, string key, string settingDefault) { if (config == null) @@ -163,6 +174,17 @@ public static TimeSpan GetTimeSpanProperty(this IProviderConfiguration config, s string s; return config.Properties.TryGetValue(key, out s) ? TimeSpan.Parse(s) : settingDefault; } + + public static bool TryGetTimeSpanProperty(this IProviderConfiguration config, string key, out TimeSpan setting) + { + if (config == null) + { + throw new ArgumentNullException("config"); + } + string s; + setting = TimeSpan.Zero; + return config.Properties.TryGetValue(key, out s) ? TimeSpan.TryParse(s, out setting) : false; + } } /// diff --git a/src/OrleansServiceBus/OrleansServiceBus.csproj b/src/OrleansServiceBus/OrleansServiceBus.csproj index eea46d620be..98380fe7e9e 100644 --- a/src/OrleansServiceBus/OrleansServiceBus.csproj +++ b/src/OrleansServiceBus/OrleansServiceBus.csproj @@ -47,6 +47,10 @@ + + + + diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs new file mode 100644 index 00000000000..45cbcb782a6 --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs @@ -0,0 +1,67 @@ +using Orleans.Runtime; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace Orleans.ServiceBus.Providers +{ + /// + /// Aggregated cache pressure monitor + /// + public class AggregatedCachePressureMonitor : List, ICachePressureMonitor + { + private bool isUnderPressure; + private Logger logger; + + /// + /// Constructor + /// + /// + public AggregatedCachePressureMonitor(Logger logger) + { + this.isUnderPressure = false; + this.logger = logger.GetSubLogger(this.GetType().Name); + } + + /// + /// Record cache pressure to every monitor in this aggregated cache monitor group + /// + /// + public void RecordCachePressureContribution(double cachePressureContribution) + { + this.ForEach(monitor => + { + monitor.RecordCachePressureContribution(cachePressureContribution); + }); + } + + /// + /// Add one monitor to this aggregated cache monitor group + /// + /// + public void AddCachePressureMonitor(ICachePressureMonitor monitor) + { + this.Add(monitor); + } + + /// + /// If any mornitor in this aggregated cache monitor group is under pressure, then return true + /// + /// + /// + public bool IsUnderPressure(DateTime utcNow) + { + bool underPressure = this.Any(monitor => monitor.IsUnderPressure(utcNow)); + if (this.isUnderPressure != underPressure) + { + this.isUnderPressure = underPressure; + logger.Info(this.isUnderPressure + ? $"Ingesting messages too fast. Throttling message reading." + : $"Message ingestion is healthy."); + } + return underPressure; + } + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs new file mode 100644 index 00000000000..a4116a02e23 --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs @@ -0,0 +1,93 @@ +using Orleans.Runtime; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace Orleans.ServiceBus.Providers +{ + /// + /// Cache pressure monitor whose back pressure algorithm is based on averaging pressure value + /// over all pressure contribution + /// + public class AveragingCachePressureMonitor : ICachePressureMonitor + { + /// + /// Default flow control threshold + /// + public static readonly double DefaultThreshold = 1 / 3; + private static readonly TimeSpan checkPeriod = TimeSpan.FromSeconds(2); + private readonly Logger logger; + + private double accumulatedCachePressure; + private double cachePressureContributionCount; + private DateTime nextCheckedTime; + private bool isUnderPressure; + private double flowControlThreshold; + + /// + /// Constructor + /// + /// + public AveragingCachePressureMonitor(Logger logger) + :this(DefaultThreshold, logger) + { } + + /// + /// Contructor + /// + /// + /// + public AveragingCachePressureMonitor(double flowControlThreshold, Logger logger) + { + this.flowControlThreshold = flowControlThreshold; + this.logger = logger.GetSubLogger(this.GetType().Name); + nextCheckedTime = DateTime.MinValue; + isUnderPressure = false; + } + + public void RecordCachePressureContribution(double cachePressureContribution) + { + // Weight unhealthy contributions thrice as much as healthy ones. + // This is a crude compensation for the fact that healthy consumers wil consume more often than unhealthy ones. + double weight = cachePressureContribution < flowControlThreshold ? 1.0 : 3.0; + accumulatedCachePressure += cachePressureContribution * weight; + cachePressureContributionCount += weight; + } + + public bool IsUnderPressure(DateTime utcNow) + { + if (nextCheckedTime < utcNow) + { + CalculatePressure(); + nextCheckedTime = utcNow + checkPeriod; + } + return isUnderPressure; + } + + private void CalculatePressure() + { + // if we don't have any contributions, don't change status + if (cachePressureContributionCount < 0.5) + { + // after 5 checks with no contributions, check anyway + cachePressureContributionCount += 0.1; + return; + } + + double pressure = accumulatedCachePressure / cachePressureContributionCount; + bool wasUnderPressure = isUnderPressure; + isUnderPressure = pressure > flowControlThreshold; + // If we changed state, log + if (isUnderPressure != wasUnderPressure) + { + logger.Verbose(isUnderPressure + ? $"Ingesting messages too fast. Throttling message reading. AccumulatedCachePressure: {accumulatedCachePressure}, Contributions: {cachePressureContributionCount}, AverageCachePressure: {pressure}, Threshold: {flowControlThreshold}" + : $"Message ingestion is healthy. AccumulatedCachePressure: {accumulatedCachePressure}, Contributions: {cachePressureContributionCount}, AverageCachePressure: {pressure}, Threshold: {flowControlThreshold}"); + } + cachePressureContributionCount = 0.0; + accumulatedCachePressure = 0.0; + } + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/ICachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/ICachePressureMonitor.cs new file mode 100644 index 00000000000..ab7f4ffe034 --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/ICachePressureMonitor.cs @@ -0,0 +1,28 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace Orleans.ServiceBus.Providers +{ + /// + /// Cache pressure monitor records pressure contribution to the cache, and determine if the cache is under pressure based on its + /// back pressure algorithm + /// + public interface ICachePressureMonitor + { + /// + /// Record cache pressure contribution to the monitor + /// + /// + void RecordCachePressureContribution(double cachePressureContribution); + + /// + /// Determine if the monitor is under pressure + /// + /// + /// + bool IsUnderPressure(DateTime utcNow); + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs new file mode 100644 index 00000000000..a56c900a46e --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -0,0 +1,114 @@ +using Orleans.Runtime; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace Orleans.ServiceBus.Providers +{ + /// + /// Pressure monitor which is in favor of the slow consumer in the cache + /// + public class SlowConsumingPressureMonitor : ICachePressureMonitor + { + /// + /// DefaultPressureWindowSize + /// + public static TimeSpan DefaultPressureWindowSize = TimeSpan.FromMinutes(1); + private const double DefaultFlowControlThreshold = 0.5; + + /// + /// PressureWindowSize + /// + public TimeSpan PressureWindowSize { get; set; } + /// + /// FlowControlThreshold + /// + public double FlowControlThreshold { get; set; } + + private readonly Logger logger; + private double biggestPressureInCurrentWindow; + private DateTime nextCheckedTime; + private bool wasUnderPressure; + + /// + /// Constructor + /// + /// + public SlowConsumingPressureMonitor(Logger logger) + : this(DefaultFlowControlThreshold, DefaultPressureWindowSize, logger) + { } + + /// + /// Constructor + /// + /// + /// + public SlowConsumingPressureMonitor(TimeSpan pressureWindowSize, Logger logger) + : this(DefaultFlowControlThreshold, pressureWindowSize, logger) + { + } + + /// + /// Constructor + /// + /// + /// + public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) + : this(flowControlThreshold, DefaultPressureWindowSize, logger) + { + } + + /// + /// Constructor + /// + /// + /// + /// + public SlowConsumingPressureMonitor(double flowControlThreshold, TimeSpan pressureWindowSzie, Logger logger) + { + this.FlowControlThreshold = flowControlThreshold; + this.logger = logger.GetSubLogger(this.GetType().Name); + this.nextCheckedTime = DateTime.MinValue; + this.biggestPressureInCurrentWindow = 0; + this.wasUnderPressure = false; + this.PressureWindowSize = pressureWindowSzie; + } + + public void RecordCachePressureContribution(double cachePressureContribution) + { + if (cachePressureContribution > this.biggestPressureInCurrentWindow) + biggestPressureInCurrentWindow = cachePressureContribution; + } + + public bool IsUnderPressure(DateTime utcNow) + { + //if any pressure contribution in current period is bigger than flowControlThreshold + //we see the cache is under pressure + bool underPressure = this.biggestPressureInCurrentWindow > this.FlowControlThreshold; + + if (underPressure && !this.wasUnderPressure) + { + //if under pressure, extend the nextCheckedTime, make sure wasUnderPressure is true for a whole window + this.wasUnderPressure = underPressure; + this.nextCheckedTime = utcNow + this.PressureWindowSize; + logger.Verbose($"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentWindow}, Threshold: {FlowControlThreshold}"); + this.biggestPressureInCurrentWindow = 0; + } + + if (this.nextCheckedTime < utcNow) + { + //at the end of each check period, reset biggestPressureInCurrentPeriod + this.nextCheckedTime = utcNow + this.PressureWindowSize; + this.biggestPressureInCurrentWindow = 0; + //if at the end of the window, pressure clears out, log + if(this.wasUnderPressure && !underPressure) + logger.Verbose($"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentWindow}, Threshold: {FlowControlThreshold}"); + this.wasUnderPressure = underPressure; + } + + return this.wasUnderPressure; + } + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs index 2cef4b6d5dd..281171bcbd8 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs @@ -143,7 +143,10 @@ public virtual void Init(IProviderConfiguration providerCfg, string providerName { var bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(1 << 20)); var timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); - CacheFactory = (partition,checkpointer,cacheLogger) => new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); + CacheFactory = (partition, checkpointer, cacheLogger) => + { + return CreateCacheFactory(partition, checkpointer, cacheLogger, bufferPool, timePurge); + }; } if (StreamFailureHandlerFactory == null) @@ -257,6 +260,29 @@ private EventHubAdapterReceiver GetOrCreateReceiver(QueueId queueId) return receivers.GetOrAdd(queueId, q => MakeReceiver(queueId)); } + private IEventHubQueueCache CreateCacheFactory(string partition, IStreamQueueCheckpointer checkpointer, Logger cacheLogger, + FixedSizeObjectPool bufferPool, TimePurgePredicate timePurge) + { + var cache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); + if (adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.HasValue) + { + var avgMonitor = new AveragingCachePressureMonitor(adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.Value, cacheLogger); +cache.AddCachePressureMonitor(avgMonitor); + } + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue + || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + { + + var slowConsumeMonitor = new SlowConsumingPressureMonitor(cacheLogger); + if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue) + slowConsumeMonitor.PressureWindowSize = adapterSettings.SlowConsumingMonitorPressureWindowSize.Value; +cache.AddCachePressureMonitor(slowConsumeMonitor); + } + return cache; + } + private EventHubAdapterReceiver MakeReceiver(QueueId queueId) { var config = new EventHubPartitionSettings diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs index 1bafc666201..cb91a16f08d 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs @@ -27,7 +27,7 @@ public abstract class EventHubQueueCache : IEventHubQueueCache /// Underlying message cache implementation /// protected readonly PooledQueueCache cache; - private readonly AveragingCachePressureMonitor cachePressureMonitor; + private readonly AggregatedCachePressureMonitor cachePressureMonitor; /// /// Logic used to store queue position @@ -38,20 +38,27 @@ public abstract class EventHubQueueCache : IEventHubQueueCache /// Construct EventHub queue cache. /// /// Default max number of items that can be added to the cache between purge calls. - /// percentage of unprocesses cache that triggers flow control /// Logic used to store queue position. /// Performs data transforms appropriate for the various types of queue data. /// Compares cached data /// - protected EventHubQueueCache(int defaultMaxAddCount, double flowControlThreshold, IStreamQueueCheckpointer checkpointer, ICacheDataAdapter cacheDataAdapter, ICacheDataComparer comparer, Logger logger) + protected EventHubQueueCache(int defaultMaxAddCount, IStreamQueueCheckpointer checkpointer, ICacheDataAdapter cacheDataAdapter, ICacheDataComparer comparer, Logger logger) { this.defaultMaxAddCount = defaultMaxAddCount; Checkpointer = checkpointer; cache = new PooledQueueCache(cacheDataAdapter, comparer, logger); cacheDataAdapter.PurgeAction = cache.Purge; cache.OnPurged = OnPurge; + this.cachePressureMonitor = new AggregatedCachePressureMonitor(logger); + } - cachePressureMonitor = new AveragingCachePressureMonitor(flowControlThreshold, logger); + /// + /// Add cache pressure monitor to the cache's back pressure algorithm + /// + /// + public void AddCachePressureMonitor(ICachePressureMonitor monitor) + { + this.cachePressureMonitor.AddCachePressureMonitor(monitor); } /// @@ -148,76 +155,11 @@ public bool TryGetNextMessage(object cursorObj, out IBatchContainer message) } - internal class AveragingCachePressureMonitor - { - private static readonly TimeSpan checkPeriod = TimeSpan.FromSeconds(2); - private readonly Logger logger; - - private double accumulatedCachePressure; - private double cachePressureContributionCount; - private DateTime nextCheckedTime; - private bool isUnderPressure; - private double flowControlThreshold; - - public AveragingCachePressureMonitor(double flowControlThreshold, Logger logger) - { - this.flowControlThreshold = flowControlThreshold; - this.logger = logger.GetSubLogger("flowcontrol", "-"); - nextCheckedTime = DateTime.MinValue; - isUnderPressure = false; - } - - public void RecordCachePressureContribution(double cachePressureContribution) - { - // Weight unhealthy contributions thrice as much as healthy ones. - // This is a crude compensation for the fact that healthy consumers wil consume more often than unhealthy ones. - double weight = cachePressureContribution < flowControlThreshold ? 1.0 : 3.0; - accumulatedCachePressure += cachePressureContribution * weight; - cachePressureContributionCount += weight; - } - - public bool IsUnderPressure(DateTime utcNow) - { - if (nextCheckedTime < utcNow) - { - CalculatePressure(); - nextCheckedTime = utcNow + checkPeriod; - } - return isUnderPressure; - } - - private void CalculatePressure() - { - // if we don't have any contributions, don't change status - if (cachePressureContributionCount < 0.5) - { - // after 5 checks with no contributions, check anyway - cachePressureContributionCount += 0.1; - return; - } - - double pressure = accumulatedCachePressure / cachePressureContributionCount; - bool wasUnderPressure = isUnderPressure; - isUnderPressure = pressure > flowControlThreshold; - // If we changed state, log - if (isUnderPressure != wasUnderPressure) - { - logger.Info(isUnderPressure - ? $"Ingesting messages too fast. Throttling message reading. AccumulatedCachePressure: {accumulatedCachePressure}, Contributions: {cachePressureContributionCount}, AverageCachePressure: {pressure}, Threshold: {flowControlThreshold}" - : $"Message ingestion is healthy. AccumulatedCachePressure: {accumulatedCachePressure}, Contributions: {cachePressureContributionCount}, AverageCachePressure: {pressure}, Threshold: {flowControlThreshold}"); - } - cachePressureContributionCount = 0.0; - accumulatedCachePressure = 0.0; - } - } - /// /// Message cache that stores EventData as a CachedEventHubMessage in a pooled message cache /// public class EventHubQueueCache : EventHubQueueCache - { - private const double DefaultThreashold = 1.0 / 3.0; - + { private readonly Logger log; /// @@ -241,24 +183,23 @@ public EventHubQueueCache(IStreamQueueCheckpointer checkpointer, IObject /// compares stream information to cached data /// cache logger public EventHubQueueCache(IStreamQueueCheckpointer checkpointer, ICacheDataAdapter cacheDataAdapter, ICacheDataComparer comparer, Logger logger) - : base(EventHubAdapterReceiver.MaxMessagesPerRead, DefaultThreashold, checkpointer, cacheDataAdapter, comparer, logger) + : base(EventHubAdapterReceiver.MaxMessagesPerRead, checkpointer, cacheDataAdapter, comparer, logger) { - log = logger.GetSubLogger("messagecache", "-"); + log = logger.GetSubLogger(this.GetType().Name); } /// /// Construct cache given a custom data adapter. /// /// Max number of message that can be added to cache from single read - /// percentage of unprocesses cache that triggers flow control /// queue checkpoint writer /// adapts queue data to cache /// compares stream information to cached data /// cache logger - public EventHubQueueCache(int defaultMaxAddCount, double flowControlThreshold, IStreamQueueCheckpointer checkpointer, ICacheDataAdapter cacheDataAdapter, ICacheDataComparer comparer, Logger logger) - : base(defaultMaxAddCount, flowControlThreshold, checkpointer, cacheDataAdapter, comparer, logger) + public EventHubQueueCache(int defaultMaxAddCount, IStreamQueueCheckpointer checkpointer, ICacheDataAdapter cacheDataAdapter, ICacheDataComparer comparer, Logger logger) + : base(defaultMaxAddCount, checkpointer, cacheDataAdapter, comparer, logger) { - log = logger.GetSubLogger("messagecache", "-"); + log = logger.GetSubLogger(this.GetType().Name); } /// @@ -311,7 +252,6 @@ protected override bool TryCalculateCachePressureContribution(StreamSequenceToke IEventHubPartitionLocation location = (IEventHubPartitionLocation) token; double cacheSize = cache.Newest.Value.SequenceNumber - cache.Oldest.Value.SequenceNumber; long distanceFromNewestMessage = cache.Newest.Value.SequenceNumber - location.SequenceNumber; - // pressure is the ratio of the distance from the front of the cache to the cachePressureContribution = distanceFromNewestMessage/cacheSize; diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs index 70b6685d18e..01f22a431ea 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs @@ -16,6 +16,37 @@ public class EventHubStreamProviderSettings /// public string StreamProviderName { get; } + /// + /// SlowConsumingMonitorFlowControlThresholdName + /// + public const string SlowConsumingMonitorFlowControlThresholdName = nameof(SlowConsumingMonitorFlowControlThreshold); + + /// + /// SlowConsumingPressureMonitorConfig + /// + public double? SlowConsumingMonitorFlowControlThreshold { get; set; } + + /// + /// SlowConsumingMonitorPressureWindowSizeName + /// + public const string SlowConsumingMonitorPressureWindowSizeName = nameof(SlowConsumingMonitorPressureWindowSize); + + /// + /// SlowConsumingMonitorPressureWindowSize + /// + public TimeSpan? SlowConsumingMonitorPressureWindowSize { get; set; } + + /// + /// AveragingCachePressureMonitorFlowControlThresholdName + /// + public const string AveragingCachePressureMonitorFlowControlThresholdName = nameof(AveragingCachePressureMonitorFlowControlThreshold); + + /// + /// AveragingCachePressureMonitorFlowControlThreshold, AveragingCachePressureMonitor is turn on by default. + /// User can turn it off by setting this value to null + /// + public double? AveragingCachePressureMonitorFlowControlThreshold = AveragingCachePressureMonitor.DefaultThreshold; + /// /// EventHubSettingsType setting name. /// @@ -110,6 +141,18 @@ public void WriteProperties(Dictionary properties) { properties.Add(DataMaxAgeInCacheName, DataMaxAgeInCache.ToString()); } + if (AveragingCachePressureMonitorFlowControlThreshold.HasValue) + { + properties.Add(AveragingCachePressureMonitorFlowControlThresholdName, AveragingCachePressureMonitorFlowControlThreshold.ToString()); + } + if (SlowConsumingMonitorPressureWindowSize.HasValue) + { + properties.Add(SlowConsumingMonitorPressureWindowSizeName, SlowConsumingMonitorPressureWindowSize.ToString()); + } + if (SlowConsumingMonitorFlowControlThreshold.HasValue) + { + properties.Add(SlowConsumingMonitorFlowControlThresholdName, SlowConsumingMonitorFlowControlThreshold.ToString()); + } } /// @@ -127,6 +170,20 @@ public void PopulateFromProviderConfig(IProviderConfiguration providerConfigurat CacheSizeMb = providerConfiguration.GetIntProperty(CacheSizeMbName, DefaultCacheSizeMb); DataMinTimeInCache = providerConfiguration.GetTimeSpanProperty(DataMinTimeInCacheName, DefaultDataMinTimeInCache); DataMaxAgeInCache = providerConfiguration.GetTimeSpanProperty(DataMaxAgeInCacheName, DefaultDataMaxAgeInCache); + double flowControlThreshold = 0; + if (providerConfiguration.TryGetDoubleProperty(SlowConsumingMonitorFlowControlThresholdName, out flowControlThreshold)) + { + this.SlowConsumingMonitorFlowControlThreshold = flowControlThreshold; + } + TimeSpan pressureWindowSize = TimeSpan.Zero; + if (providerConfiguration.TryGetTimeSpanProperty(SlowConsumingMonitorPressureWindowSizeName, out pressureWindowSize)) + { + this.SlowConsumingMonitorPressureWindowSize = pressureWindowSize; + } + if (providerConfiguration.TryGetDoubleProperty(AveragingCachePressureMonitorFlowControlThresholdName, out flowControlThreshold)) + { + this.AveragingCachePressureMonitorFlowControlThreshold = flowControlThreshold; + } } /// diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs index 8fece915d2c..507abf3de16 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs @@ -36,5 +36,11 @@ public interface IEventHubQueueCache : IQueueFlowController, IDisposable /// /// bool TryGetNextMessage(object cursorObj, out IBatchContainer message); + + /// + /// Add cache pressure monitor to the cache's back pressure algorithm + /// + /// + void AddCachePressureMonitor(ICachePressureMonitor monitor); } } diff --git a/test/ServiceBus.Tests/ServiceBus.Tests.csproj b/test/ServiceBus.Tests/ServiceBus.Tests.csproj index a5260e2bf26..c0ed1dfa350 100644 --- a/test/ServiceBus.Tests/ServiceBus.Tests.csproj +++ b/test/ServiceBus.Tests/ServiceBus.Tests.csproj @@ -127,12 +127,15 @@ + + + diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs new file mode 100644 index 00000000000..365dde9921a --- /dev/null +++ b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs @@ -0,0 +1,194 @@ +using Microsoft.WindowsAzure.Storage.Table; +using Orleans; +using Orleans.AzureUtils; +using Orleans.Providers.Streams.Common; +using Orleans.Runtime; +using Orleans.Runtime.Configuration; +using Orleans.ServiceBus.Providers; +using Orleans.Storage; +using Orleans.Streams; +using Orleans.TestingHost; +using Orleans.TestingHost.Utils; +using ServiceBus.Tests.TestStreamProviders; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using TestExtensions; +using UnitTests.GrainInterfaces; +using UnitTests.Grains.ProgrammaticSubscribe; +using Xunit; + +namespace ServiceBus.Tests.SlowConsumingTests +{ + [Collection(TestEnvironmentFixture.DefaultCollection)] + public class EHSlowConsumingTests : OrleansTestingBase, IClassFixture + { + private const string StreamProviderName = "EventHubStreamProvider"; + private const string StreamNamespace = "EHTestsNamespace"; + private const string EHPath = "ehorleanstest"; + private const string EHConsumerGroup = "orleansnightly"; + private const string EHCheckpointTable = "ehcheckpoint"; + private static readonly string CheckpointNamespace = Guid.NewGuid().ToString(); + private static readonly TimeSpan monitorPressureWindowSize = TimeSpan.FromSeconds(3); + private static readonly TimeSpan timeout = TimeSpan.FromSeconds(30); + private const double flowControlThredhold = 0.6; + public static readonly EventHubStreamProviderSettings ProviderSettings = + new EventHubStreamProviderSettings(StreamProviderName); + + private static readonly Lazy EventHubConfig = new Lazy(() => + new EventHubSettings( + TestDefaultConfiguration.EventHubConnectionString, + EHConsumerGroup, EHPath)); + + private static readonly EventHubCheckpointerSettings CheckpointerSettings = + new EventHubCheckpointerSettings(TestDefaultConfiguration.DataConnectionString, + EHCheckpointTable, CheckpointNamespace, TimeSpan.FromSeconds(1)); + + + private readonly Fixture fixture; + + public class Fixture : BaseTestClusterFixture + { + protected override TestCluster CreateTestCluster() + { + var options = new TestClusterOptions(2); + ProviderSettings.SlowConsumingMonitorPressureWindowSize = monitorPressureWindowSize; + ProviderSettings.SlowConsumingMonitorFlowControlThreshold = flowControlThredhold; + ProviderSettings.AveragingCachePressureMonitorFlowControlThreshold = null; + AdjustClusterConfiguration(options.ClusterConfiguration); + return new TestCluster(options); + } + + private bool isSkippable; + protected override void CheckPreconditionsOrThrow() + { + base.CheckPreconditionsOrThrow(); + if (string.IsNullOrWhiteSpace(TestDefaultConfiguration.EventHubConnectionString) || + string.IsNullOrWhiteSpace(TestDefaultConfiguration.DataConnectionString)) + { + this.isSkippable = true; + throw new SkipException("EventHubConnectionString or DataConnectionString is not set up"); + } + } + + public override void Dispose() + { + base.Dispose(); + if (!isSkippable) + { + var dataManager = new AzureTableDataManager(CheckpointerSettings.TableName, CheckpointerSettings.DataConnectionString); + dataManager.InitTableAsync().Wait(); + dataManager.ClearTableAsync().Wait(); + } + } + + private static void AdjustClusterConfiguration(ClusterConfiguration config) + { + var settings = new Dictionary(); + // get initial settings from configs + ProviderSettings.WriteProperties(settings); + EventHubConfig.Value.WriteProperties(settings); + CheckpointerSettings.WriteProperties(settings); + + // add queue balancer setting + settings.Add(PersistentStreamProviderConfig.QUEUE_BALANCER_TYPE, StreamQueueBalancerType.DynamicClusterConfigDeploymentBalancer.ToString()); + + // register stream provider + config.Globals.RegisterStreamProvider(StreamProviderName, settings); + config.Globals.RegisterStorageProvider("PubSubStore"); + } + } + + public EHSlowConsumingTests(Fixture fixture) + { + this.fixture = fixture; + fixture.EnsurePreconditionsMet(); + } + + [SkippableFact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("Functional")] + public async Task EHSlowConsuming_ShouldFavorSlowConsumer() + { + var streamId = new FullStreamIdentity(Guid.NewGuid(), StreamNamespace, StreamProviderName); + //set up one slow consumer grain + var slowConsumer = this.fixture.GrainFactory.GetGrain(Guid.NewGuid()); + await slowConsumer.BecomeConsumer(streamId.Guid, StreamNamespace, StreamProviderName); + + //set up 30 healthy consumer grain to show how much we favor slow consumer + int healthyConsumerCount = 30; + var healthyConsumers = await SetUpHealthyConsumerGrain(this.fixture.GrainFactory, streamId.Guid, StreamNamespace, StreamProviderName, healthyConsumerCount); + + //set up producer and start producing + var producer = this.fixture.GrainFactory.GetGrain(Guid.NewGuid()); + await producer.BecomeProducer(streamId.Guid, StreamNamespace, StreamProviderName); + await producer.StartPeriodicProducing(); + + //since there's an extreme slow consumer, so the back pressure algorithm should be triggered + await TestingUtils.WaitUntilAsync(lastTry => AssertCacheBackPressureTriggered(true, lastTry), timeout); + + //make slow consumer stop consuming + await slowConsumer.StopConsuming(); + + //slowConsumer stopped consuming, back pressure algorithm should be cleared in next check period. + await Task.Delay(monitorPressureWindowSize); + await TestingUtils.WaitUntilAsync(lastTry => AssertCacheBackPressureTriggered(false, lastTry), timeout); + + //clean up test + await producer.StopPeriodicProducing(); + await StopHealthyConsumerGrainComing(healthyConsumers); + } + + private async Task> SetUpHealthyConsumerGrain(IGrainFactory GrainFactory, Guid streamId, string streamNameSpace, string streamProvider, int grainCount) + { + List grains = new List(); + List tasks = new List(); + while (grainCount > 0) + { + var consumer = GrainFactory.GetGrain(Guid.NewGuid()); + grains.Add(consumer); + tasks.Add(consumer.BecomeConsumer(streamId, streamNameSpace, streamProvider)); + grainCount--; + } + await Task.WhenAll(tasks); + return grains; + } + + private async Task StopHealthyConsumerGrainComing(List grains) + { + List tasks = new List(); + foreach (var grain in grains) + { + tasks.Add(grain.StopConsuming()); + } + await Task.WhenAll(tasks); + } + + private async Task AssertCacheBackPressureTriggered(bool expectedResult, bool assertIsTrue) + { + if (assertIsTrue) + { + bool actualResult = await IsBackPressureTriggered(); + Assert.True(expectedResult == actualResult, $"Back pressure algorithm should be triggered? expected: {expectedResult}, actual: {actualResult}"); + return true; + } + else + { + return (await IsBackPressureTriggered()) == expectedResult; + } + } + + private async Task IsBackPressureTriggered() + { + IManagementGrain mgmtGrain = this.fixture.HostedCluster.GrainFactory.GetGrain(0); + object[] replies = await mgmtGrain.SendControlCommandToProvider(typeof(EHStreamProviderWithCreatedCacheList).FullName, + StreamProviderName, EHStreamProviderWithCreatedCacheList.AdapterFactory.IsCacheBackPressureTriggeredCommand, null); + foreach (var re in replies) + { + if ((bool)re) + return true; + } + return false; + } + } +} diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs new file mode 100644 index 00000000000..843b5932ce4 --- /dev/null +++ b/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs @@ -0,0 +1,69 @@ +using Orleans.Providers; +using Orleans.Runtime.Configuration; +using Orleans.ServiceBus.Providers; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using Xunit; + +namespace NonSilo.Tests.StreamingTests +{ + + /// + /// EventHubAdapterFactory populate EventHubStreamProviderSettings from IProviderConfiguration. + /// So this test suit tests that EventHubStreamProviderSettings will be populated back as the same as before + /// it is written into ProviderConfiguration. + /// + public class EventHubStreamProviderSettingsTests + { + private static string StreamProviderName = "EHStreamProvider"; + [Fact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("BVT")] + public void DefaultSetting_Write_Into_ProviderConfiguration_PopulateBack() + { + var expectedSetting = new EventHubStreamProviderSettings(StreamProviderName); + AssertSettingEqual_After_WriteInto_ProviderConfiguration_AndPopulateBack(expectedSetting); + } + + [Fact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("BVT")] + public void SettingWithSlowConsumingMonitorSetUp_Write_Into_ProviderConfiguration_PopulateBack() + { + var expectedSetting = new EventHubStreamProviderSettings(StreamProviderName); + expectedSetting.SlowConsumingMonitorPressureWindowSize = TimeSpan.FromMinutes(2); + expectedSetting.SlowConsumingMonitorFlowControlThreshold = 1 / 3; + AssertSettingEqual_After_WriteInto_ProviderConfiguration_AndPopulateBack(expectedSetting); + } + + [Fact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("BVT")] + public void SettingWithAvgConsumingMonitorSetUp_Write_Into_ProviderConfiguration_PopulateBack() + { + var expectedSetting = new EventHubStreamProviderSettings(StreamProviderName); + expectedSetting.AveragingCachePressureMonitorFlowControlThreshold = 1 / 10; + AssertSettingEqual_After_WriteInto_ProviderConfiguration_AndPopulateBack(expectedSetting); + } + + private void AssertSettingEqual_After_WriteInto_ProviderConfiguration_AndPopulateBack(EventHubStreamProviderSettings expectedSetting) + { + var properties = new Dictionary(); + expectedSetting.WriteProperties(properties); + var config = new ProviderConfiguration(properties, typeof(EventHubStreamProvider).FullName, StreamProviderName); + + var actualSettings = new EventHubStreamProviderSettings(StreamProviderName); + actualSettings.PopulateFromProviderConfig(config); + AssertEqual(expectedSetting, actualSettings); + } + private void AssertEqual(EventHubStreamProviderSettings expectedSettings, EventHubStreamProviderSettings actualSettings) + { + Assert.Equal(expectedSettings.StreamProviderName, actualSettings.StreamProviderName); + Assert.Equal(expectedSettings.SlowConsumingMonitorFlowControlThreshold, actualSettings.SlowConsumingMonitorFlowControlThreshold); + Assert.Equal(expectedSettings.SlowConsumingMonitorPressureWindowSize, actualSettings.SlowConsumingMonitorPressureWindowSize); + Assert.Equal(expectedSettings.AveragingCachePressureMonitorFlowControlThreshold, actualSettings.AveragingCachePressureMonitorFlowControlThreshold); + Assert.Equal(expectedSettings.EventHubSettingsType, actualSettings.EventHubSettingsType); + Assert.Equal(expectedSettings.CheckpointerSettingsType, actualSettings.CheckpointerSettingsType); + Assert.Equal(expectedSettings.CacheSizeMb, actualSettings.CacheSizeMb); + Assert.Equal(expectedSettings.DataMinTimeInCache, actualSettings.DataMinTimeInCache); + Assert.Equal(expectedSettings.DataMaxAgeInCache, actualSettings.DataMaxAgeInCache); + } + } +} diff --git a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs new file mode 100644 index 00000000000..a34eab2d834 --- /dev/null +++ b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs @@ -0,0 +1,78 @@ +using Orleans.Providers; +using Orleans.Providers.Streams.Common; +using Orleans.Runtime; +using Orleans.ServiceBus.Providers; +using Orleans.Streams; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace ServiceBus.Tests.TestStreamProviders +{ + class EHStreamProviderWithCreatedCacheList : PersistentStreamProvider + { + public class AdapterFactory : EventHubAdapterFactory, IControllable + { + private List createdCaches; + private FixedSizeObjectPool bufferPool; + private TimePurgePredicate timePurge; + private static int defaultMaxAddCount = 10; + public AdapterFactory() + { + createdCaches = new List(); + CacheFactory = CreateQueueCache; + } + + private IEventHubQueueCache CreateQueueCache(string partition, IStreamQueueCheckpointer checkpointer, Logger cacheLogger) + { + // the same code block as with EventHubAdapterFactory to define default CacheFactory + // except for at the end we put the created cache in CreatedCaches list + if(this.bufferPool == null) + this.bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(1 << 20)); + if(this.timePurge == null) + this.timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); + + //ser defaultMaxAddCount to 10 so TryCalculateCachePressureContribution will start to calculate real contribution shortly. + var cache = new EventHubQueueCache(defaultMaxAddCount, checkpointer, new EventHubDataAdapter(this.SerializationManager, bufferPool, timePurge), EventHubDataComparer.Instance, cacheLogger); + if (adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.HasValue) + { + var avgMonitor = new AveragingCachePressureMonitor(adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.Value, cacheLogger); + cache.AddCachePressureMonitor(avgMonitor); + } + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue + || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + { + + var slowConsumeMonitor = new SlowConsumingPressureMonitor(cacheLogger); + if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue) + slowConsumeMonitor.PressureWindowSize = adapterSettings.SlowConsumingMonitorPressureWindowSize.Value; + cache.AddCachePressureMonitor(slowConsumeMonitor); + } + this.createdCaches.Add(cache); + return cache; + } + + public static int IsCacheBackPressureTriggeredCommand = (int)PersistentStreamProviderCommand.AdapterFactoryCommandStartRange + 3; + /// + /// Only command expecting: determine whether back pressure algorithm on any of the created caches + /// is triggered. + /// + /// + /// + /// + public Task ExecuteCommand(int command, object arg) + { + foreach (var cache in this.createdCaches) + { + if (cache.GetMaxAddCount() == 0) + return Task.FromResult(true); + } + return Task.FromResult(false); + } + } + } +} diff --git a/test/TestGrainInterfaces/SlowConsumingGrains/ISlowConsumingGrain.cs b/test/TestGrainInterfaces/SlowConsumingGrains/ISlowConsumingGrain.cs new file mode 100644 index 00000000000..be64c78dd9a --- /dev/null +++ b/test/TestGrainInterfaces/SlowConsumingGrains/ISlowConsumingGrain.cs @@ -0,0 +1,19 @@ +using Orleans; +using Orleans.Streams; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace UnitTests.GrainInterfaces +{ + public interface ISlowConsumingGrain : IGrainWithGuidKey + { + Task BecomeConsumer(Guid streamId, string streamNamespace, string providerToUse); + + Task StopConsuming(); + + Task GetNumberConsumed(); + } +} diff --git a/test/TestGrainInterfaces/TestGrainInterfaces.csproj b/test/TestGrainInterfaces/TestGrainInterfaces.csproj index 8627e09ebea..5b60660657d 100644 --- a/test/TestGrainInterfaces/TestGrainInterfaces.csproj +++ b/test/TestGrainInterfaces/TestGrainInterfaces.csproj @@ -132,6 +132,7 @@ + diff --git a/test/TestGrains/ProgrammaticSubscribe/Passive_ConsumerGrain.cs b/test/TestGrains/ProgrammaticSubscribe/Passive_ConsumerGrain.cs index cd6e199dfa3..87f7ae9d5f7 100644 --- a/test/TestGrains/ProgrammaticSubscribe/Passive_ConsumerGrain.cs +++ b/test/TestGrains/ProgrammaticSubscribe/Passive_ConsumerGrain.cs @@ -144,7 +144,7 @@ public class ConsumerObserver : IAsyncObserver internal ConsumerObserver(Logger logger) { this.NumConsumed = 0; - this.logger = logger; + this.logger = logger.GetSubLogger(this.GetType().Name); } public Task OnNextAsync(T item, StreamSequenceToken token = null) diff --git a/test/TestGrains/SlowConsumingGrains/SlowConsumingGrain.cs b/test/TestGrains/SlowConsumingGrains/SlowConsumingGrain.cs new file mode 100644 index 00000000000..4235e9c160a --- /dev/null +++ b/test/TestGrains/SlowConsumingGrains/SlowConsumingGrain.cs @@ -0,0 +1,91 @@ +using Orleans; +using Orleans.Runtime; +using Orleans.Streams; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using UnitTests.GrainInterfaces; + +namespace UnitTests.Grains +{ + /// + /// SlowConsumingGrain keep asking to rewind to the first item it received to mimic slow consuming behavior + /// + public class SlowConsumingGrain : Grain, ISlowConsumingGrain + { + private Logger logger; + public SlowObserver ConsumerObserver { get; private set; } + public StreamSubscriptionHandle ConsumerHandle { get; set; } + public override Task OnActivateAsync() + { + logger = base.GetLogger(this.GetType().Name + base.IdentityString); + logger.Info("OnActivateAsync"); + ConsumerHandle = null; + return TaskDone.Done; + } + + public Task GetNumberConsumed() + { + return Task.FromResult(this.ConsumerObserver.NumConsumed); + } + + public async Task BecomeConsumer(Guid streamId, string streamNamespace, string providerToUse) + { + logger.Info("BecomeConsumer"); + ConsumerObserver = new SlowObserver(this, logger); + IStreamProvider streamProvider = base.GetStreamProvider(providerToUse); + var consumer = streamProvider.GetStream(streamId, streamNamespace); + ConsumerHandle = await consumer.SubscribeAsync(ConsumerObserver); + } + + public async Task StopConsuming() + { + logger.Info("StopConsuming"); + if (ConsumerHandle != null) + { + await ConsumerHandle.UnsubscribeAsync(); + ConsumerHandle = null; + } + } + } + + /// + /// SlowObserver keep rewind to the first item it received, to mimic slow consuming behavior + /// + /// + public class SlowObserver : IAsyncObserver + { + public int NumConsumed { get; private set; } + private Logger logger; + private SlowConsumingGrain slowConsumingGrain; + internal SlowObserver(SlowConsumingGrain grain, Logger logger) + { + NumConsumed = 0; + this.slowConsumingGrain = grain; + this.logger = logger.GetSubLogger(this.GetType().Name); + } + + public async Task OnNextAsync(T item, StreamSequenceToken token = null) + { + NumConsumed++; + // slow consumer keep asking for the first item it received to mimic slow consuming behavior + this.slowConsumingGrain.ConsumerHandle = await this.slowConsumingGrain.ConsumerHandle.ResumeAsync(this.slowConsumingGrain.ConsumerObserver, token); + this.logger.Info($"Consumer {this.GetHashCode()} OnNextAsync() received item {item.ToString()}, with NumConsumed {NumConsumed}"); + } + + public Task OnCompletedAsync() + { + this.logger.Info($"Consumer {this.GetHashCode()} OnCompletedAsync()"); + return TaskDone.Done; + } + + public Task OnErrorAsync(Exception ex) + { + this.logger.Info($"Consumer {this.GetHashCode()} OnErrorAsync({ex})"); + return TaskDone.Done; + } + } + +} diff --git a/test/TestGrains/TestGrains.csproj b/test/TestGrains/TestGrains.csproj index c23065d2c33..8a3c107ad8e 100644 --- a/test/TestGrains/TestGrains.csproj +++ b/test/TestGrains/TestGrains.csproj @@ -118,6 +118,7 @@ +