-
Notifications
You must be signed in to change notification settings - Fork 2k
/
EventHubQueueMapper.cs
62 lines (59 loc) · 2.35 KB
/
EventHubQueueMapper.cs
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
using System;
using System.Collections.Generic;
using System.Globalization;
using System.Linq;
using Orleans.Streams;
using Orleans.Configuration;
namespace Orleans.ServiceBus.Providers
{
/// <summary>
/// Queue mapper that tracks which EventHub partition was mapped to which queueId
/// </summary>
public class EventHubQueueMapper : HashRingBasedStreamQueueMapper, IEventHubQueueMapper
{
private readonly Dictionary<QueueId, string> partitionDictionary = new Dictionary<QueueId, string>();
private static HashRingStreamQueueMapperOptions GetHashRingStreamQueueMapperOptions(string[] partitionIds)
{
var options = new HashRingStreamQueueMapperOptions();
options.TotalQueueCount = partitionIds.Length;
return options;
}
/// <summary>
/// Queue mapper that tracks which EventHub partition was mapped to which queueId
/// </summary>
/// <param name="partitionIds">List of EventHubPartitions</param>
/// <param name="queueNamePrefix">Prefix for queueIds. Must be unique per stream provider</param>
public EventHubQueueMapper(string[] partitionIds, string queueNamePrefix)
: base(GetHashRingStreamQueueMapperOptions(partitionIds), queueNamePrefix)
{
QueueId[] queues = GetAllQueues().ToArray();
if (queues.Length != partitionIds.Length)
{
throw new ArgumentOutOfRangeException(nameof(partitionIds), "partitions and Queues do not line up");
}
for (int i = 0; i < queues.Length; i++)
{
partitionDictionary.Add(queues[i], partitionIds[i]);
}
}
/// <summary>
/// Gets the EventHub partition by QueueId
/// </summary>
/// <param name="queue"></param>
/// <returns></returns>
public string QueueToPartition(QueueId queue)
{
if (queue == null)
{
throw new ArgumentNullException(nameof(queue));
}
string partitionId;
if (!partitionDictionary.TryGetValue(queue, out partitionId))
{
throw new ArgumentOutOfRangeException(string.Format(CultureInfo.InvariantCulture, "queue {0}", queue.ToStringWithHashCode()));
}
return partitionId;
}
}
}