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