From 378d76bbe525c7a2484ad1d0cfc274cd3aaf8446 Mon Sep 17 00:00:00 2001 From: xiazen Date: Wed, 22 Mar 2017 17:15:02 -0700 Subject: [PATCH 1/8] add SlowConsumingPressureMonitor --- .../Streams/EventHub/EventHubQueueCache.cs | 67 +++++++++++++++++-- .../Streams/EventHub/ICachePressureMonitor.cs | 45 +++++++++++++ .../Streams/EventHub/IEventHubQueueCache.cs | 7 ++ 3 files changed, 115 insertions(+), 4 deletions(-) create mode 100644 src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs index 1bafc666201..a3c54540b08 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs @@ -9,6 +9,7 @@ using Orleans.Runtime; using Orleans.Serialization; using Orleans.Streams; +using OrleansServiceBus.Providers.Streams.EventHub; namespace Orleans.ServiceBus.Providers { @@ -27,7 +28,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 @@ -50,8 +51,18 @@ protected EventHubQueueCache(int defaultMaxAddCount, double flowControlThreshold cache = new PooledQueueCache(cacheDataAdapter, comparer, logger); cacheDataAdapter.PurgeAction = cache.Purge; cache.OnPurged = OnPurge; + + var avgCachePressureMonitor = new AveragingCachePressureMonitor(flowControlThreshold, logger); + this.cachePressureMonitor = new AggregatedCachePressureMonitor() { avgCachePressureMonitor }; + } - 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,7 +159,55 @@ public bool TryGetNextMessage(object cursorObj, out IBatchContainer message) } - internal class AveragingCachePressureMonitor + public class SlowConsumingPressureMonitor : ICachePressureMonitor + { + private readonly TimeSpan checkPeriod = TimeSpan.FromMinutes(1); + private readonly Logger logger; + + private double biggestPressureInCurrentPeriod; + private DateTime nextCheckedTime; + private double flowControlThreshold; + private bool isUnderPressure; + + public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) + { + this.flowControlThreshold = flowControlThreshold; + this.logger = logger.GetSubLogger("flowcontrol-slow-consumer-pressure", "-"); + this.nextCheckedTime = DateTime.MinValue; + this.biggestPressureInCurrentPeriod = 0; + this.isUnderPressure = false; + } + + public void RecordCachePressureContribution(double cachePressureContribution) + { + if (cachePressureContribution > biggestPressureInCurrentPeriod) + biggestPressureInCurrentPeriod = 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.biggestPressureInCurrentPeriod > this.flowControlThreshold; + if (this.isUnderPressure != underPressure) + { + this.isUnderPressure = underPressure; + logger.Info(this.isUnderPressure + ? $"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {flowControlThreshold}" + : $"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {flowControlThreshold}"); + } + + if (nextCheckedTime < utcNow) + { + //at the end of each check period, reset biggestPressureInCurrentPeriod + this.nextCheckedTime = utcNow + this.checkPeriod; + this.biggestPressureInCurrentPeriod = 0; + } + return underPressure; + } + } + + internal class AveragingCachePressureMonitor : ICachePressureMonitor { private static readonly TimeSpan checkPeriod = TimeSpan.FromSeconds(2); private readonly Logger logger; @@ -162,7 +221,7 @@ internal class AveragingCachePressureMonitor public AveragingCachePressureMonitor(double flowControlThreshold, Logger logger) { this.flowControlThreshold = flowControlThreshold; - this.logger = logger.GetSubLogger("flowcontrol", "-"); + this.logger = logger.GetSubLogger("flowcontrol-averaging-cache-pressure", "-"); nextCheckedTime = DateTime.MinValue; isUnderPressure = false; } diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs new file mode 100644 index 00000000000..d174ad82bd9 --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs @@ -0,0 +1,45 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace OrleansServiceBus.Providers.Streams.EventHub +{ + public interface ICachePressureMonitor + { + void RecordCachePressureContribution(double cachePressureContribution); + + bool IsUnderPressure(DateTime utcNow); + } + + internal class AggregatedCachePressureMonitor : List, ICachePressureMonitor + { + public void RecordCachePressureContribution(double cachePressureContribution) + { + this.ForEach(monitor => + { + monitor.RecordCachePressureContribution(cachePressureContribution); + }); + } + + public void AddCachePressureMonitor(ICachePressureMonitor monitor) + { + this.Add(monitor); + } + + public bool IsUnderPressure(DateTime utcNow) + { + bool isUnderPressure = false; + //if any mornitor in this monitor list is under pressure, then return true + this.ForEach(monitor => + { + if (monitor.IsUnderPressure(utcNow)) + { + isUnderPressure = true; + } + }); + return isUnderPressure; + } + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs index 8fece915d2c..1405d9e5de9 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs @@ -7,6 +7,7 @@ #endif using Orleans.Providers.Streams.Common; using Orleans.Streams; +using OrleansServiceBus.Providers.Streams.EventHub; namespace Orleans.ServiceBus.Providers { @@ -36,5 +37,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); } } From 300db7c4d10654ce408885b2844bf84ee0f5b907 Mon Sep 17 00:00:00 2001 From: xiazen Date: Wed, 22 Mar 2017 18:14:05 -0700 Subject: [PATCH 2/8] add one option to configure EventHubQeueuCache with the SlowConsumingPressureMonitor --- .../OrleansServiceBus.csproj | 1 + test/ServiceBus.Tests/ServiceBus.Tests.csproj | 1 + ...viderWithSlowConsumingPressureDetecting.cs | 33 +++++++++++++++++++ 3 files changed, 35 insertions(+) create mode 100644 test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs diff --git a/src/OrleansServiceBus/OrleansServiceBus.csproj b/src/OrleansServiceBus/OrleansServiceBus.csproj index eea46d620be..fc594f8b7cf 100644 --- a/src/OrleansServiceBus/OrleansServiceBus.csproj +++ b/src/OrleansServiceBus/OrleansServiceBus.csproj @@ -63,6 +63,7 @@ + diff --git a/test/ServiceBus.Tests/ServiceBus.Tests.csproj b/test/ServiceBus.Tests/ServiceBus.Tests.csproj index a5260e2bf26..28a34d4f51c 100644 --- a/test/ServiceBus.Tests/ServiceBus.Tests.csproj +++ b/test/ServiceBus.Tests/ServiceBus.Tests.csproj @@ -133,6 +133,7 @@ + diff --git a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs new file mode 100644 index 00000000000..e1f8c576869 --- /dev/null +++ b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs @@ -0,0 +1,33 @@ +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 +{ + public class EHStreamProviderWithSlowConsumingPressureDetecting : PersistentStreamProvider + { + public class AdapterFactory : EventHubAdapterFactory + { + public AdapterFactory() + { + CacheFactory = CreateQueueCache; + } + + private IEventHubQueueCache CreateQueueCache(string partition, IStreamQueueCheckpointer checkpointer, Logger log) + { + var bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(1 << 20)); + var timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); + var eventhubQueeuCache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, log, this.SerializationManager); + var slowConsumingPressureMonitor = new SlowConsumingPressureMonitor(0.5, log); + eventhubQueeuCache.AddCachePressureMonitor(slowConsumingPressureMonitor); + return eventhubQueeuCache; + } + } + } +} From 93723674bf3657f3d8a7eb59fd0674861634466f Mon Sep 17 00:00:00 2001 From: xiazen Date: Thu, 23 Mar 2017 12:46:56 -0700 Subject: [PATCH 3/8] add configure option through EventHubProviderSettings --- .../EventHub/EventHubAdapterFactory.cs | 15 +++- .../Streams/EventHub/EventHubQueueCache.cs | 9 ++- .../EventHubStreamProviderSettings.cs | 2 + .../Streaming/EHSlowConsumingTests.cs | 75 +++++++++++++++++++ ...viderWithSlowConsumingPressureDetecting.cs | 5 +- 5 files changed, 102 insertions(+), 4 deletions(-) create mode 100644 test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs index 2cef4b6d5dd..a840146ad3a 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs @@ -143,7 +143,20 @@ 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); + if (adapterSettings.SlowConsumingMonitorThreshold > 0) + { + CacheFactory = (partition, checkpointer, cacheLogger) => + { var cache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); + var monitor = new SlowConsumingPressureMonitor(adapterSettings.SlowConsumingMonitorThreshold, log); + cache.AddCachePressureMonitor(monitor); + return cache; + }; + } + else + { + CacheFactory = (partition, checkpointer, cacheLogger) => new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); + } + } if (StreamFailureHandlerFactory == null) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs index a3c54540b08..8553282370e 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs @@ -161,7 +161,8 @@ public bool TryGetNextMessage(object cursorObj, out IBatchContainer message) public class SlowConsumingPressureMonitor : ICachePressureMonitor { - private readonly TimeSpan checkPeriod = TimeSpan.FromMinutes(1); + private static TimeSpan defaultCheckPeriod = TimeSpan.FromMinutes(1); + private readonly TimeSpan checkPeriod; private readonly Logger logger; private double biggestPressureInCurrentPeriod; @@ -170,12 +171,18 @@ public class SlowConsumingPressureMonitor : ICachePressureMonitor private bool isUnderPressure; public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) + :this(flowControlThreshold, logger, defaultCheckPeriod) + { + } + + public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger, TimeSpan checkPeriod) { this.flowControlThreshold = flowControlThreshold; this.logger = logger.GetSubLogger("flowcontrol-slow-consumer-pressure", "-"); this.nextCheckedTime = DateTime.MinValue; this.biggestPressureInCurrentPeriod = 0; this.isUnderPressure = false; + this.checkPeriod = checkPeriod; } public void RecordCachePressureContribution(double cachePressureContribution) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs index 70b6685d18e..f16ce5772ac 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs @@ -16,6 +16,8 @@ public class EventHubStreamProviderSettings /// public string StreamProviderName { get; } + public double SlowConsumingMonitorThreshold { get; set; } + /// /// EventHubSettingsType setting name. /// diff --git a/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs new file mode 100644 index 00000000000..4fb66160c7e --- /dev/null +++ b/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs @@ -0,0 +1,75 @@ +using Orleans.Runtime.Configuration; +using Orleans.ServiceBus.Providers; +using Orleans.Storage; +using Orleans.Streams; +using Orleans.TestingHost; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using TestExtensions; +using Xunit; + +namespace ServiceBus.Tests.Streaming +{ + [TestCategory("EventHub"), TestCategory("Streaming")] + [Collection(TestEnvironmentFixture.DefaultCollection)] + class EHSlowConsumingTests : OrleansTestingBase, IClassFixture + { + private const string StreamProviderName = "EventHubStreamProvider"; + private const string StreamNamespace = "EHSlowConsumingTestsNamespace"; + private const string EHPath = "ehorleanstest"; + private const string EHConsumerGroup = "orleansnightly"; + private const string EHCheckpointTable = "ehcheckpoint"; + private static readonly string CheckpointNamespace = Guid.NewGuid().ToString(); + + 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); + AdjustClusterConfiguration(options.ClusterConfiguration); + return new TestCluster(options); + } + + private static void AdjustClusterConfiguration(ClusterConfiguration config) + { + var settings = new Dictionary(); + + //configure slow consuming monitor threshhold + ProviderSettings.SlowConsumingMonitorThreshold = 0.5; + // 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; + } + } +} diff --git a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs index e1f8c576869..b2a891eae61 100644 --- a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs +++ b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs @@ -21,10 +21,11 @@ public AdapterFactory() private IEventHubQueueCache CreateQueueCache(string partition, IStreamQueueCheckpointer checkpointer, Logger log) { - var bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(1 << 20)); + var blockSize = 1 << 20; + var bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(blockSize)); var timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); var eventhubQueeuCache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, log, this.SerializationManager); - var slowConsumingPressureMonitor = new SlowConsumingPressureMonitor(0.5, log); + var slowConsumingPressureMonitor = new SlowConsumingPressureMonitor(0.5, log, TimeSpan.FromMinutes(1)); eventhubQueeuCache.AddCachePressureMonitor(slowConsumingPressureMonitor); return eventhubQueeuCache; } From 6908027c8c8f5e3b5726fee83e726903efee15e2 Mon Sep 17 00:00:00 2001 From: xiazen Date: Mon, 27 Mar 2017 14:03:44 -0700 Subject: [PATCH 4/8] PR feedback and add testing --- src/Orleans/Providers/IOrleansProvider.cs | 22 +++ .../OrleansServiceBus.csproj | 5 +- .../AggregatedCachePressureMonitor.cs} | 29 ++- .../AveragingCachePressureMonitor.cs | 93 +++++++++ .../ICachePressureMonitor.cs | 28 +++ .../SlowConsumingPressureMonitor.cs | 104 ++++++++++ .../EventHub/EventHubAdapterFactory.cs | 33 ++-- .../Streams/EventHub/EventHubQueueCache.cs | 142 +------------- .../EventHubStreamProviderSettings.cs | 57 +++++- .../Streams/EventHub/IEventHubQueueCache.cs | 1 - test/ServiceBus.Tests/ServiceBus.Tests.csproj | 4 +- .../EHSlowConsumingTests.cs | 178 ++++++++++++++++++ .../EventHubStreamProviderSettingsTests.cs | 69 +++++++ .../EHStreamProviderWithCreatedCacheList.cs | 78 ++++++++ ...viderWithSlowConsumingPressureDetecting.cs | 34 ---- .../ISlowConsumingGrain.cs | 19 ++ .../TestGrainInterfaces.csproj | 1 + .../Passive_ConsumerGrain.cs | 2 +- .../SlowConsumingGrains/SlowConsumingGrain.cs | 91 +++++++++ test/TestGrains/TestGrains.csproj | 1 + 20 files changed, 795 insertions(+), 196 deletions(-) rename src/OrleansServiceBus/Providers/Streams/EventHub/{ICachePressureMonitor.cs => CachePressureMonitors/AggregatedCachePressureMonitor.cs} (51%) create mode 100644 src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs create mode 100644 src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/ICachePressureMonitor.cs create mode 100644 src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs create mode 100644 test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs create mode 100644 test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs create mode 100644 test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs delete mode 100644 test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs create mode 100644 test/TestGrainInterfaces/SlowConsumingGrains/ISlowConsumingGrain.cs create mode 100644 test/TestGrains/SlowConsumingGrains/SlowConsumingGrain.cs 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 fc594f8b7cf..98380fe7e9e 100644 --- a/src/OrleansServiceBus/OrleansServiceBus.csproj +++ b/src/OrleansServiceBus/OrleansServiceBus.csproj @@ -47,6 +47,10 @@ + + + + @@ -63,7 +67,6 @@ - diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs similarity index 51% rename from src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs rename to src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs index d174ad82bd9..90789001980 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/ICachePressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs @@ -4,17 +4,17 @@ using System.Text; using System.Threading.Tasks; -namespace OrleansServiceBus.Providers.Streams.EventHub +namespace Orleans.ServiceBus.Providers { - public interface ICachePressureMonitor - { - void RecordCachePressureContribution(double cachePressureContribution); - - bool IsUnderPressure(DateTime utcNow); - } - - internal class AggregatedCachePressureMonitor : List, ICachePressureMonitor + /// + /// Aggregated cache pressure monitor + /// + public class AggregatedCachePressureMonitor : List, ICachePressureMonitor { + /// + /// Record cache pressure to every monitor in this aggregated cache monitor group + /// + /// public void RecordCachePressureContribution(double cachePressureContribution) { this.ForEach(monitor => @@ -23,15 +23,24 @@ public void RecordCachePressureContribution(double 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 isUnderPressure = false; - //if any mornitor in this monitor list is under pressure, then return true + this.ForEach(monitor => { if (monitor.IsUnderPressure(utcNow)) 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..937324d0c05 --- /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 DefaultThreashold = 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(DefaultThreashold, 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.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; + } + } +} 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..1def45ce039 --- /dev/null +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -0,0 +1,104 @@ +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 + { + private static TimeSpan DefaultCheckPeriod = TimeSpan.FromMinutes(1); + private const double DefaultFlowControlThreshold = 0.5; + + /// + /// CheckPeriod + /// + public TimeSpan CheckPeriod { get; set; } + /// + /// FlowControlThreshold + /// + public double FlowControlThreshold { get; set; } + + private readonly Logger logger; + private double biggestPressureInCurrentPeriod; + private DateTime nextCheckedTime; + private bool isUnderPressure; + + /// + /// Constructor + /// + /// + public SlowConsumingPressureMonitor(Logger logger) + : this(DefaultFlowControlThreshold, DefaultCheckPeriod, logger) + { } + + /// + /// Constructor + /// + /// + /// + public SlowConsumingPressureMonitor(TimeSpan checkPeriod, Logger logger) + : this(DefaultFlowControlThreshold, checkPeriod, logger) + { + } + + /// + /// Constructor + /// + /// + /// + public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) + : this(flowControlThreshold, DefaultCheckPeriod, logger) + { + } + + /// + /// Constructor + /// + /// + /// + /// + public SlowConsumingPressureMonitor(double flowControlThreshold, TimeSpan checkPeriod, Logger logger) + { + this.FlowControlThreshold = flowControlThreshold; + this.logger = logger.GetSubLogger(this.GetType().Name); + this.nextCheckedTime = DateTime.MinValue; + this.biggestPressureInCurrentPeriod = 0; + this.isUnderPressure = false; + this.CheckPeriod = checkPeriod; + } + + public void RecordCachePressureContribution(double cachePressureContribution) + { + if (cachePressureContribution > biggestPressureInCurrentPeriod) + biggestPressureInCurrentPeriod = 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.biggestPressureInCurrentPeriod > this.FlowControlThreshold; + if (this.isUnderPressure != underPressure) + { + this.isUnderPressure = underPressure; + logger.Info(this.isUnderPressure + ? $"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}" + : $"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}"); + } + + if (nextCheckedTime < utcNow) + { + //at the end of each check period, reset biggestPressureInCurrentPeriod + this.nextCheckedTime = utcNow + this.CheckPeriod; + this.biggestPressureInCurrentPeriod = 0; + } + return underPressure; + } + } +} diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs index a840146ad3a..1bbc68941bd 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs @@ -143,20 +143,27 @@ 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); - if (adapterSettings.SlowConsumingMonitorThreshold > 0) + CacheFactory = (partition, checkpointer, cacheLogger) => { - CacheFactory = (partition, checkpointer, cacheLogger) => - { var cache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); - var monitor = new SlowConsumingPressureMonitor(adapterSettings.SlowConsumingMonitorThreshold, log); - cache.AddCachePressureMonitor(monitor); - return cache; - }; - } - else - { - CacheFactory = (partition, checkpointer, cacheLogger) => new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); - } - + var cache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); + if (adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.HasValue) + { + var avgMonitor = new AveragingCachePressureMonitor(adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.Value, log); + cache.AddCachePressureMonitor(avgMonitor); + } + if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue + || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + { + + var slowConsumeMonitor = new SlowConsumingPressureMonitor(log); + if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; + if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue) + slowConsumeMonitor.CheckPeriod = adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.Value; + cache.AddCachePressureMonitor(slowConsumeMonitor); + } + return cache; + }; } if (StreamFailureHandlerFactory == null) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs index 8553282370e..b504a0013cb 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs @@ -9,7 +9,6 @@ using Orleans.Runtime; using Orleans.Serialization; using Orleans.Streams; -using OrleansServiceBus.Providers.Streams.EventHub; namespace Orleans.ServiceBus.Providers { @@ -39,21 +38,18 @@ 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; - - var avgCachePressureMonitor = new AveragingCachePressureMonitor(flowControlThreshold, logger); - this.cachePressureMonitor = new AggregatedCachePressureMonitor() { avgCachePressureMonitor }; + this.cachePressureMonitor = new AggregatedCachePressureMonitor(); } /// @@ -159,131 +155,11 @@ public bool TryGetNextMessage(object cursorObj, out IBatchContainer message) } - public class SlowConsumingPressureMonitor : ICachePressureMonitor - { - private static TimeSpan defaultCheckPeriod = TimeSpan.FromMinutes(1); - private readonly TimeSpan checkPeriod; - private readonly Logger logger; - - private double biggestPressureInCurrentPeriod; - private DateTime nextCheckedTime; - private double flowControlThreshold; - private bool isUnderPressure; - - public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) - :this(flowControlThreshold, logger, defaultCheckPeriod) - { - } - - public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger, TimeSpan checkPeriod) - { - this.flowControlThreshold = flowControlThreshold; - this.logger = logger.GetSubLogger("flowcontrol-slow-consumer-pressure", "-"); - this.nextCheckedTime = DateTime.MinValue; - this.biggestPressureInCurrentPeriod = 0; - this.isUnderPressure = false; - this.checkPeriod = checkPeriod; - } - - public void RecordCachePressureContribution(double cachePressureContribution) - { - if (cachePressureContribution > biggestPressureInCurrentPeriod) - biggestPressureInCurrentPeriod = 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.biggestPressureInCurrentPeriod > this.flowControlThreshold; - if (this.isUnderPressure != underPressure) - { - this.isUnderPressure = underPressure; - logger.Info(this.isUnderPressure - ? $"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {flowControlThreshold}" - : $"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {flowControlThreshold}"); - } - - if (nextCheckedTime < utcNow) - { - //at the end of each check period, reset biggestPressureInCurrentPeriod - this.nextCheckedTime = utcNow + this.checkPeriod; - this.biggestPressureInCurrentPeriod = 0; - } - return underPressure; - } - } - - internal class AveragingCachePressureMonitor : ICachePressureMonitor - { - 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-averaging-cache-pressure", "-"); - 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; /// @@ -307,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); } /// @@ -377,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 f16ce5772ac..5084e61e79c 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs @@ -16,7 +16,36 @@ public class EventHubStreamProviderSettings /// public string StreamProviderName { get; } - public double SlowConsumingMonitorThreshold { get; set; } + /// + /// SlowConsumingMonitorFlowControlThresholdName + /// + public const string SlowConsumingMonitorFlowControlThresholdName = "SlowConsumingMonitorFlowControlThreshold"; + + /// + /// SlowConsumingPressureMonitorConfig + /// + public double? SlowConsumingMonitorFlowControlThreshold { get; set; } + + /// + /// SlowConsumingMonitorFlowControlCheckperiodName + /// + public const string SlowConsumingMonitorFlowControlCheckperiodName = "SlowConsumingMonitorFlowControlCheckperiod"; + + /// + /// SlowConsumingMonitorFlowControlCheckperiod + /// + public TimeSpan? SlowConsumingMonitorFlowControlCheckperiod { get; set; } + + /// + /// AveragingCachePressureMonitorFlowControlThresholdName + /// + public const string AveragingCachePressureMonitorFlowControlThresholdName = "AveragingCachePressureMonitorFlowControlThreshold"; + + /// + /// AveragingCachePressureMonitorFlowControlThreshold, AveragingCachePressureMonitor is turn on by default. + /// User can turn it off by setting this value to null + /// + public double? AveragingCachePressureMonitorFlowControlThreshold = AveragingCachePressureMonitor.DefaultThreashold; /// /// EventHubSettingsType setting name. @@ -112,6 +141,18 @@ public void WriteProperties(Dictionary properties) { properties.Add(DataMaxAgeInCacheName, DataMaxAgeInCache.ToString()); } + if (AveragingCachePressureMonitorFlowControlThreshold.HasValue) + { + properties.Add(AveragingCachePressureMonitorFlowControlThresholdName, AveragingCachePressureMonitorFlowControlThreshold.ToString()); + } + if (SlowConsumingMonitorFlowControlCheckperiod.HasValue) + { + properties.Add(SlowConsumingMonitorFlowControlCheckperiodName, SlowConsumingMonitorFlowControlCheckperiod.ToString()); + } + if (SlowConsumingMonitorFlowControlThreshold.HasValue) + { + properties.Add(SlowConsumingMonitorFlowControlThresholdName, SlowConsumingMonitorFlowControlThreshold.ToString()); + } } /// @@ -129,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 checkPeriod = TimeSpan.Zero; + if (providerConfiguration.TryGetTimeSpanProperty(SlowConsumingMonitorFlowControlCheckperiodName, out checkPeriod)) + { + this.SlowConsumingMonitorFlowControlCheckperiod = checkPeriod; + } + 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 1405d9e5de9..507abf3de16 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/IEventHubQueueCache.cs @@ -7,7 +7,6 @@ #endif using Orleans.Providers.Streams.Common; using Orleans.Streams; -using OrleansServiceBus.Providers.Streams.EventHub; namespace Orleans.ServiceBus.Providers { diff --git a/test/ServiceBus.Tests/ServiceBus.Tests.csproj b/test/ServiceBus.Tests/ServiceBus.Tests.csproj index 28a34d4f51c..c0ed1dfa350 100644 --- a/test/ServiceBus.Tests/ServiceBus.Tests.csproj +++ b/test/ServiceBus.Tests/ServiceBus.Tests.csproj @@ -127,13 +127,15 @@ + + - + diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs new file mode 100644 index 00000000000..56524854017 --- /dev/null +++ b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs @@ -0,0 +1,178 @@ +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 monitorCheckperiod = TimeSpan.FromSeconds(30); + 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.SlowConsumingMonitorFlowControlCheckperiod = monitorCheckperiod; + ProviderSettings.SlowConsumingMonitorFlowControlThreshold = flowControlThredhold; + ProviderSettings.AveragingCachePressureMonitorFlowControlThreshold = null; + AdjustClusterConfiguration(options.ClusterConfiguration); + return new TestCluster(options); + } + + public override void Dispose() + { + base.Dispose(); + 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; + } + + [Fact, 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(monitorCheckperiod); + 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..6313a1fbb7d --- /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.SlowConsumingMonitorFlowControlCheckperiod = 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.SlowConsumingMonitorFlowControlCheckperiod, actualSettings.SlowConsumingMonitorFlowControlCheckperiod); + 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..e203ae8f3fd --- /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.SlowConsumingMonitorFlowControlCheckperiod.HasValue + || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + { + + var slowConsumeMonitor = new SlowConsumingPressureMonitor(cacheLogger); + if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) + slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; + if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue) + slowConsumeMonitor.CheckPeriod = adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.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/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs deleted file mode 100644 index b2a891eae61..00000000000 --- a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithSlowConsumingPressureDetecting.cs +++ /dev/null @@ -1,34 +0,0 @@ -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 -{ - public class EHStreamProviderWithSlowConsumingPressureDetecting : PersistentStreamProvider - { - public class AdapterFactory : EventHubAdapterFactory - { - public AdapterFactory() - { - CacheFactory = CreateQueueCache; - } - - private IEventHubQueueCache CreateQueueCache(string partition, IStreamQueueCheckpointer checkpointer, Logger log) - { - var blockSize = 1 << 20; - var bufferPool = new FixedSizeObjectPool(adapterSettings.CacheSizeMb, () => new FixedSizeBuffer(blockSize)); - var timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); - var eventhubQueeuCache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, log, this.SerializationManager); - var slowConsumingPressureMonitor = new SlowConsumingPressureMonitor(0.5, log, TimeSpan.FromMinutes(1)); - eventhubQueeuCache.AddCachePressureMonitor(slowConsumingPressureMonitor); - return eventhubQueeuCache; - } - } - } -} 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 @@ + From 6562a19243bd5c42282f7fcf257ed6a7cb8b13eb Mon Sep 17 00:00:00 2001 From: Xiao Zeng Date: Sun, 2 Apr 2017 11:54:28 -0700 Subject: [PATCH 5/8] PR feedback --- .../AggregatedCachePressureMonitor.cs | 11 +-- .../AveragingCachePressureMonitor.cs | 4 +- .../SlowConsumingPressureMonitor.cs | 24 +++--- .../EventHub/EventHubAdapterFactory.cs | 42 ++++++----- .../EventHubStreamProviderSettings.cs | 24 +++--- .../EHSlowConsumingTests.cs | 30 ++++++-- .../EventHubStreamProviderSettingsTests.cs | 4 +- .../Streaming/EHSlowConsumingTests.cs | 75 ------------------- .../EHStreamProviderWithCreatedCacheList.cs | 6 +- 9 files changed, 79 insertions(+), 141 deletions(-) delete mode 100644 test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs index 90789001980..c903a3c6c62 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs @@ -39,16 +39,7 @@ public void AddCachePressureMonitor(ICachePressureMonitor monitor) /// public bool IsUnderPressure(DateTime utcNow) { - bool isUnderPressure = false; - - this.ForEach(monitor => - { - if (monitor.IsUnderPressure(utcNow)) - { - isUnderPressure = true; - } - }); - return isUnderPressure; + return this.Any(monitor => monitor.IsUnderPressure(utcNow)); } } } diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs index 937324d0c05..86b2b650a8d 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs @@ -16,7 +16,7 @@ public class AveragingCachePressureMonitor : ICachePressureMonitor /// /// Default flow control threshold /// - public static readonly double DefaultThreashold = 1 / 3; + public static readonly double DefaultThreshold = 1 / 3; private static readonly TimeSpan checkPeriod = TimeSpan.FromSeconds(2); private readonly Logger logger; @@ -31,7 +31,7 @@ public class AveragingCachePressureMonitor : ICachePressureMonitor /// /// public AveragingCachePressureMonitor(Logger logger) - :this(DefaultThreashold, logger) + :this(DefaultThreshold, logger) { } /// diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs index 1def45ce039..4747820c2e5 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -12,13 +12,13 @@ namespace Orleans.ServiceBus.Providers /// public class SlowConsumingPressureMonitor : ICachePressureMonitor { - private static TimeSpan DefaultCheckPeriod = TimeSpan.FromMinutes(1); + private static TimeSpan DefaultPressureWindowSize = TimeSpan.FromMinutes(1); private const double DefaultFlowControlThreshold = 0.5; /// - /// CheckPeriod + /// PressureWindowSize /// - public TimeSpan CheckPeriod { get; set; } + public TimeSpan PressureWindowSize { get; set; } /// /// FlowControlThreshold /// @@ -34,16 +34,16 @@ public class SlowConsumingPressureMonitor : ICachePressureMonitor /// /// public SlowConsumingPressureMonitor(Logger logger) - : this(DefaultFlowControlThreshold, DefaultCheckPeriod, logger) + : this(DefaultFlowControlThreshold, DefaultPressureWindowSize, logger) { } /// /// Constructor /// - /// + /// /// - public SlowConsumingPressureMonitor(TimeSpan checkPeriod, Logger logger) - : this(DefaultFlowControlThreshold, checkPeriod, logger) + public SlowConsumingPressureMonitor(TimeSpan pressureWindowSize, Logger logger) + : this(DefaultFlowControlThreshold, pressureWindowSize, logger) { } @@ -53,7 +53,7 @@ public SlowConsumingPressureMonitor(TimeSpan checkPeriod, Logger logger) /// /// public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) - : this(flowControlThreshold, DefaultCheckPeriod, logger) + : this(flowControlThreshold, DefaultPressureWindowSize, logger) { } @@ -61,16 +61,16 @@ public SlowConsumingPressureMonitor(double flowControlThreshold, Logger logger) /// Constructor /// /// - /// + /// /// - public SlowConsumingPressureMonitor(double flowControlThreshold, TimeSpan checkPeriod, Logger logger) + public SlowConsumingPressureMonitor(double flowControlThreshold, TimeSpan pressureWindowSzie, Logger logger) { this.FlowControlThreshold = flowControlThreshold; this.logger = logger.GetSubLogger(this.GetType().Name); this.nextCheckedTime = DateTime.MinValue; this.biggestPressureInCurrentPeriod = 0; this.isUnderPressure = false; - this.CheckPeriod = checkPeriod; + this.PressureWindowSize = pressureWindowSzie; } public void RecordCachePressureContribution(double cachePressureContribution) @@ -95,7 +95,7 @@ public bool IsUnderPressure(DateTime utcNow) if (nextCheckedTime < utcNow) { //at the end of each check period, reset biggestPressureInCurrentPeriod - this.nextCheckedTime = utcNow + this.CheckPeriod; + this.nextCheckedTime = utcNow + this.PressureWindowSize; this.biggestPressureInCurrentPeriod = 0; } return underPressure; diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs index 1bbc68941bd..281171bcbd8 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubAdapterFactory.cs @@ -145,24 +145,7 @@ public virtual void Init(IProviderConfiguration providerCfg, string providerName var timePurge = new TimePurgePredicate(adapterSettings.DataMinTimeInCache, adapterSettings.DataMaxAgeInCache); CacheFactory = (partition, checkpointer, cacheLogger) => { - var cache = new EventHubQueueCache(checkpointer, bufferPool, timePurge, cacheLogger, this.SerializationManager); - if (adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.HasValue) - { - var avgMonitor = new AveragingCachePressureMonitor(adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.Value, log); - cache.AddCachePressureMonitor(avgMonitor); - } - if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue - || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) - { - - var slowConsumeMonitor = new SlowConsumingPressureMonitor(log); - if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) - slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; - if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue) - slowConsumeMonitor.CheckPeriod = adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.Value; - cache.AddCachePressureMonitor(slowConsumeMonitor); - } - return cache; + return CreateCacheFactory(partition, checkpointer, cacheLogger, bufferPool, timePurge); }; } @@ -277,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/EventHubStreamProviderSettings.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs index 5084e61e79c..01f22a431ea 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubStreamProviderSettings.cs @@ -19,7 +19,7 @@ public class EventHubStreamProviderSettings /// /// SlowConsumingMonitorFlowControlThresholdName /// - public const string SlowConsumingMonitorFlowControlThresholdName = "SlowConsumingMonitorFlowControlThreshold"; + public const string SlowConsumingMonitorFlowControlThresholdName = nameof(SlowConsumingMonitorFlowControlThreshold); /// /// SlowConsumingPressureMonitorConfig @@ -27,25 +27,25 @@ public class EventHubStreamProviderSettings public double? SlowConsumingMonitorFlowControlThreshold { get; set; } /// - /// SlowConsumingMonitorFlowControlCheckperiodName + /// SlowConsumingMonitorPressureWindowSizeName /// - public const string SlowConsumingMonitorFlowControlCheckperiodName = "SlowConsumingMonitorFlowControlCheckperiod"; + public const string SlowConsumingMonitorPressureWindowSizeName = nameof(SlowConsumingMonitorPressureWindowSize); /// - /// SlowConsumingMonitorFlowControlCheckperiod + /// SlowConsumingMonitorPressureWindowSize /// - public TimeSpan? SlowConsumingMonitorFlowControlCheckperiod { get; set; } + public TimeSpan? SlowConsumingMonitorPressureWindowSize { get; set; } /// /// AveragingCachePressureMonitorFlowControlThresholdName /// - public const string AveragingCachePressureMonitorFlowControlThresholdName = "AveragingCachePressureMonitorFlowControlThreshold"; + 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.DefaultThreashold; + public double? AveragingCachePressureMonitorFlowControlThreshold = AveragingCachePressureMonitor.DefaultThreshold; /// /// EventHubSettingsType setting name. @@ -145,9 +145,9 @@ public void WriteProperties(Dictionary properties) { properties.Add(AveragingCachePressureMonitorFlowControlThresholdName, AveragingCachePressureMonitorFlowControlThreshold.ToString()); } - if (SlowConsumingMonitorFlowControlCheckperiod.HasValue) + if (SlowConsumingMonitorPressureWindowSize.HasValue) { - properties.Add(SlowConsumingMonitorFlowControlCheckperiodName, SlowConsumingMonitorFlowControlCheckperiod.ToString()); + properties.Add(SlowConsumingMonitorPressureWindowSizeName, SlowConsumingMonitorPressureWindowSize.ToString()); } if (SlowConsumingMonitorFlowControlThreshold.HasValue) { @@ -175,10 +175,10 @@ public void PopulateFromProviderConfig(IProviderConfiguration providerConfigurat { this.SlowConsumingMonitorFlowControlThreshold = flowControlThreshold; } - TimeSpan checkPeriod = TimeSpan.Zero; - if (providerConfiguration.TryGetTimeSpanProperty(SlowConsumingMonitorFlowControlCheckperiodName, out checkPeriod)) + TimeSpan pressureWindowSize = TimeSpan.Zero; + if (providerConfiguration.TryGetTimeSpanProperty(SlowConsumingMonitorPressureWindowSizeName, out pressureWindowSize)) { - this.SlowConsumingMonitorFlowControlCheckperiod = checkPeriod; + this.SlowConsumingMonitorPressureWindowSize = pressureWindowSize; } if (providerConfiguration.TryGetDoubleProperty(AveragingCachePressureMonitorFlowControlThresholdName, out flowControlThreshold)) { diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs index 56524854017..3833256b0cb 100644 --- a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs +++ b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs @@ -31,7 +31,7 @@ public class EHSlowConsumingTests : OrleansTestingBase, IClassFixture(CheckpointerSettings.TableName, CheckpointerSettings.DataConnectionString); - dataManager.InitTableAsync().Wait(); - dataManager.ClearTableAsync().Wait(); + if (!isSkippable) + { + var dataManager = new AzureTableDataManager(CheckpointerSettings.TableName, CheckpointerSettings.DataConnectionString); + dataManager.InitTableAsync().Wait(); + dataManager.ClearTableAsync().Wait(); + } } private static void AdjustClusterConfiguration(ClusterConfiguration config) @@ -89,9 +104,10 @@ private static void AdjustClusterConfiguration(ClusterConfiguration config) public EHSlowConsumingTests(Fixture fixture) { this.fixture = fixture; + fixture.EnsurePreconditionsMet(); } - [Fact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("Functional")] + [SkippableFact, TestCategory("EventHub"), TestCategory("Streaming"), TestCategory("Functional")] public async Task EHSlowConsuming_ShouldFavorSlowConsumer() { var streamId = new FullStreamIdentity(Guid.NewGuid(), StreamNamespace, StreamProviderName); @@ -115,7 +131,7 @@ public async Task EHSlowConsuming_ShouldFavorSlowConsumer() await slowConsumer.StopConsuming(); //slowConsumer stopped consuming, back pressure algorithm should be cleared in next check period. - await Task.Delay(monitorCheckperiod); + await Task.Delay(monitorPressureWindowSize); await TestingUtils.WaitUntilAsync(lastTry => AssertCacheBackPressureTriggered(false, lastTry), timeout); //clean up test diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs index 6313a1fbb7d..843b5932ce4 100644 --- a/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs +++ b/test/ServiceBus.Tests/SlowConsumingTests/EventHubStreamProviderSettingsTests.cs @@ -30,7 +30,7 @@ public void DefaultSetting_Write_Into_ProviderConfiguration_PopulateBack() public void SettingWithSlowConsumingMonitorSetUp_Write_Into_ProviderConfiguration_PopulateBack() { var expectedSetting = new EventHubStreamProviderSettings(StreamProviderName); - expectedSetting.SlowConsumingMonitorFlowControlCheckperiod = TimeSpan.FromMinutes(2); + expectedSetting.SlowConsumingMonitorPressureWindowSize = TimeSpan.FromMinutes(2); expectedSetting.SlowConsumingMonitorFlowControlThreshold = 1 / 3; AssertSettingEqual_After_WriteInto_ProviderConfiguration_AndPopulateBack(expectedSetting); } @@ -57,7 +57,7 @@ private void AssertEqual(EventHubStreamProviderSettings expectedSettings, EventH { Assert.Equal(expectedSettings.StreamProviderName, actualSettings.StreamProviderName); Assert.Equal(expectedSettings.SlowConsumingMonitorFlowControlThreshold, actualSettings.SlowConsumingMonitorFlowControlThreshold); - Assert.Equal(expectedSettings.SlowConsumingMonitorFlowControlCheckperiod, actualSettings.SlowConsumingMonitorFlowControlCheckperiod); + Assert.Equal(expectedSettings.SlowConsumingMonitorPressureWindowSize, actualSettings.SlowConsumingMonitorPressureWindowSize); Assert.Equal(expectedSettings.AveragingCachePressureMonitorFlowControlThreshold, actualSettings.AveragingCachePressureMonitorFlowControlThreshold); Assert.Equal(expectedSettings.EventHubSettingsType, actualSettings.EventHubSettingsType); Assert.Equal(expectedSettings.CheckpointerSettingsType, actualSettings.CheckpointerSettingsType); diff --git a/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs deleted file mode 100644 index 4fb66160c7e..00000000000 --- a/test/ServiceBus.Tests/Streaming/EHSlowConsumingTests.cs +++ /dev/null @@ -1,75 +0,0 @@ -using Orleans.Runtime.Configuration; -using Orleans.ServiceBus.Providers; -using Orleans.Storage; -using Orleans.Streams; -using Orleans.TestingHost; -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using System.Threading.Tasks; -using TestExtensions; -using Xunit; - -namespace ServiceBus.Tests.Streaming -{ - [TestCategory("EventHub"), TestCategory("Streaming")] - [Collection(TestEnvironmentFixture.DefaultCollection)] - class EHSlowConsumingTests : OrleansTestingBase, IClassFixture - { - private const string StreamProviderName = "EventHubStreamProvider"; - private const string StreamNamespace = "EHSlowConsumingTestsNamespace"; - private const string EHPath = "ehorleanstest"; - private const string EHConsumerGroup = "orleansnightly"; - private const string EHCheckpointTable = "ehcheckpoint"; - private static readonly string CheckpointNamespace = Guid.NewGuid().ToString(); - - 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); - AdjustClusterConfiguration(options.ClusterConfiguration); - return new TestCluster(options); - } - - private static void AdjustClusterConfiguration(ClusterConfiguration config) - { - var settings = new Dictionary(); - - //configure slow consuming monitor threshhold - ProviderSettings.SlowConsumingMonitorThreshold = 0.5; - // 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; - } - } -} diff --git a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs index e203ae8f3fd..a34eab2d834 100644 --- a/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs +++ b/test/ServiceBus.Tests/TestStreamProviders/EHStreamProviderWithCreatedCacheList.cs @@ -41,15 +41,15 @@ private IEventHubQueueCache CreateQueueCache(string partition, IStreamQueueCheck var avgMonitor = new AveragingCachePressureMonitor(adapterSettings.AveragingCachePressureMonitorFlowControlThreshold.Value, cacheLogger); cache.AddCachePressureMonitor(avgMonitor); } - if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue || adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) { var slowConsumeMonitor = new SlowConsumingPressureMonitor(cacheLogger); if (adapterSettings.SlowConsumingMonitorFlowControlThreshold.HasValue) slowConsumeMonitor.FlowControlThreshold = adapterSettings.SlowConsumingMonitorFlowControlThreshold.Value; - if (adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.HasValue) - slowConsumeMonitor.CheckPeriod = adapterSettings.SlowConsumingMonitorFlowControlCheckperiod.Value; + if (adapterSettings.SlowConsumingMonitorPressureWindowSize.HasValue) + slowConsumeMonitor.PressureWindowSize = adapterSettings.SlowConsumingMonitorPressureWindowSize.Value; cache.AddCachePressureMonitor(slowConsumeMonitor); } this.createdCaches.Add(cache); From 96f7de07c065969ffdb7f3823a173de2312f3d94 Mon Sep 17 00:00:00 2001 From: xiazen Date: Tue, 4 Apr 2017 13:36:00 -0700 Subject: [PATCH 6/8] PR feedback --- .../AggregatedCachePressureMonitor.cs | 26 +++++++++++++++++-- .../AveragingCachePressureMonitor.cs | 2 +- .../SlowConsumingPressureMonitor.cs | 2 +- .../Streams/EventHub/EventHubQueueCache.cs | 2 +- 4 files changed, 27 insertions(+), 5 deletions(-) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs index c903a3c6c62..45cbcb782a6 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AggregatedCachePressureMonitor.cs @@ -1,4 +1,5 @@ -using System; +using Orleans.Runtime; +using System; using System.Collections.Generic; using System.Linq; using System.Text; @@ -11,6 +12,19 @@ namespace Orleans.ServiceBus.Providers /// 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 /// @@ -39,7 +53,15 @@ public void AddCachePressureMonitor(ICachePressureMonitor monitor) /// public bool IsUnderPressure(DateTime utcNow) { - return this.Any(monitor => monitor.IsUnderPressure(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 index 86b2b650a8d..a4116a02e23 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/AveragingCachePressureMonitor.cs @@ -82,7 +82,7 @@ private void CalculatePressure() // If we changed state, log if (isUnderPressure != wasUnderPressure) { - logger.Info(isUnderPressure + 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}"); } diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs index 4747820c2e5..8c22ebbbe12 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -87,7 +87,7 @@ public bool IsUnderPressure(DateTime utcNow) if (this.isUnderPressure != underPressure) { this.isUnderPressure = underPressure; - logger.Info(this.isUnderPressure + logger.Verbose(this.isUnderPressure ? $"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}" : $"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}"); } diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs index b504a0013cb..cb91a16f08d 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/EventHubQueueCache.cs @@ -49,7 +49,7 @@ protected EventHubQueueCache(int defaultMaxAddCount, IStreamQueueCheckpointer(cacheDataAdapter, comparer, logger); cacheDataAdapter.PurgeAction = cache.Purge; cache.OnPurged = OnPurge; - this.cachePressureMonitor = new AggregatedCachePressureMonitor(); + this.cachePressureMonitor = new AggregatedCachePressureMonitor(logger); } /// From a728bca6ea6c87ad72e5c93937a06c307106149d Mon Sep 17 00:00:00 2001 From: xiazen Date: Tue, 4 Apr 2017 17:15:29 -0700 Subject: [PATCH 7/8] make slowpressuremonitor more conservative --- .../SlowConsumingPressureMonitor.cs | 37 +++++++++++-------- .../EHSlowConsumingTests.cs | 2 +- 2 files changed, 23 insertions(+), 16 deletions(-) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs index 8c22ebbbe12..4cf6f460ca6 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -25,9 +25,9 @@ public class SlowConsumingPressureMonitor : ICachePressureMonitor public double FlowControlThreshold { get; set; } private readonly Logger logger; - private double biggestPressureInCurrentPeriod; + private double biggestPressureInCurrentWindow; private DateTime nextCheckedTime; - private bool isUnderPressure; + private bool wasUnderPressure; /// /// Constructor @@ -68,37 +68,44 @@ public SlowConsumingPressureMonitor(double flowControlThreshold, TimeSpan pressu this.FlowControlThreshold = flowControlThreshold; this.logger = logger.GetSubLogger(this.GetType().Name); this.nextCheckedTime = DateTime.MinValue; - this.biggestPressureInCurrentPeriod = 0; - this.isUnderPressure = false; + this.biggestPressureInCurrentWindow = 0; + this.wasUnderPressure = false; this.PressureWindowSize = pressureWindowSzie; } public void RecordCachePressureContribution(double cachePressureContribution) { - if (cachePressureContribution > biggestPressureInCurrentPeriod) - biggestPressureInCurrentPeriod = 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.biggestPressureInCurrentPeriod > this.FlowControlThreshold; - if (this.isUnderPressure != underPressure) + bool underPressure = this.biggestPressureInCurrentWindow > this.FlowControlThreshold; + + if (underPressure && !this.wasUnderPressure) { - this.isUnderPressure = underPressure; - logger.Verbose(this.isUnderPressure - ? $"Ingesting messages too fast. Throttling message reading. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}" - : $"Message ingestion is healthy. BiggestPressureInCurrentPeriod: {biggestPressureInCurrentPeriod}, Threshold: {FlowControlThreshold}"); + //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 (nextCheckedTime < utcNow) + if (this.nextCheckedTime < utcNow) { //at the end of each check period, reset biggestPressureInCurrentPeriod this.nextCheckedTime = utcNow + this.PressureWindowSize; - this.biggestPressureInCurrentPeriod = 0; + 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 underPressure; + + return this.wasUnderPressure; } } } diff --git a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs index 3833256b0cb..365dde9921a 100644 --- a/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs +++ b/test/ServiceBus.Tests/SlowConsumingTests/EHSlowConsumingTests.cs @@ -31,7 +31,7 @@ public class EHSlowConsumingTests : OrleansTestingBase, IClassFixture Date: Thu, 6 Apr 2017 10:38:29 -0700 Subject: [PATCH 8/8] make SlowConsumerMonitor.DefaultWindowSize public --- .../CachePressureMonitors/SlowConsumingPressureMonitor.cs | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs index 4cf6f460ca6..a56c900a46e 100644 --- a/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs +++ b/src/OrleansServiceBus/Providers/Streams/EventHub/CachePressureMonitors/SlowConsumingPressureMonitor.cs @@ -12,7 +12,10 @@ namespace Orleans.ServiceBus.Providers /// public class SlowConsumingPressureMonitor : ICachePressureMonitor { - private static TimeSpan DefaultPressureWindowSize = TimeSpan.FromMinutes(1); + /// + /// DefaultPressureWindowSize + /// + public static TimeSpan DefaultPressureWindowSize = TimeSpan.FromMinutes(1); private const double DefaultFlowControlThreshold = 0.5; ///