Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
Надежная конкурентная очередь — это асинхронная, транзакционная и реплицированная очередь, которая обладает высокой степенью параллелизма для операций добавления и удаления из очереди. Он предназначен для обеспечения высокой пропускной способности и низкой задержки путем ослабления строгого порядка FIFO, предоставляемого надежной очередью, и вместо этого обеспечивает наилучший порядок.
Программные интерфейсы
| Одновременная очередь | Надежная параллельная очередь |
|---|---|
| void Enqueue(T item) | Задача EnqueueAsync(ITransaction tx, T item) |
| bool TryDequeue(out T result) | Task< ConditionalValue < T >> TryDequeueAsync(ITransaction tx) |
| int Count() | long Count() |
Сравнение с очередью с высокой надежностью
Надежная параллельная очередь предлагается в качестве альтернативы надежной очереди. Его следует использовать в тех случаях, когда строгое упорядочение FIFO не требуется, так как гарантия FIFO требует компромисса с параллелизмом. Надежная очередь использует блокировки для обеспечения порядка FIFO, при этом не более одной транзакции разрешено помещать в очередь и не более одной транзакции разрешено извлекать из очереди одновременно. В сравнении с этим, надёжная конкурентная очередь ослабляет ограничение упорядочивания и позволяет любому количеству параллельных транзакций переплетать их операции добавления и удаления из очереди. Упорядочение на основе лучших усилий предоставляется, однако относительный порядок двух значений в надежной параллельной очереди никогда не может быть гарантирован.
Надежная одновременная очередь обеспечивает более высокую пропускную способность и низкую задержку, чем Надежная очередь, когда несколько параллельных транзакций выполняют операции добавления в очередь и/или извлечения из нее.
Пример использования для ReliableConcurrentQueue — сценарий очереди сообщений . В этом сценарии один или несколько производителей сообщений создают и добавляют элементы в очередь, а один или несколько потребителей сообщений извлекают сообщения из очереди и обрабатывают их. Несколько производителей и потребителей могут работать независимо, используя одновременные транзакции для обработки очереди.
Рекомендации по использованию
- Очередь ожидает, что объекты в очереди имеют короткий срок хранения. То есть элементы не будут оставаться в очереди в течение длительного времени.
- Очередь не гарантирует строгий порядок следования FIFO.
- Очередь не считывает записи, которые сама создала. Если элемент поставлен в очередь в пределах транзакции, он не будет виден для декаера в той же транзакции.
- Очереди не изолированы друг от друга. Если элемент A извлечен из очереди в транзакции txnA, даже если txnA не зафиксирована, элемент A не будет видимым для параллельной транзакции txnB. Если txnA прерывается, A станет видимым для txnB немедленно.
- Поведение TryPeekAsync можно реализовать с помощью TryDequeueAsync и прерывания транзакции. Пример этого поведения можно найти в разделе "Шаблоны программирования".
- Подсчет нетранзакционный. Его можно использовать для получения представления о количестве элементов в очереди, но представляет точку во времени и не может полагаться.
- Не рекомендуется выполнять дорогостоящую обработку извлечённых элементов, пока транзакция активна, чтобы избежать долго выполняющихся транзакций, которые могут повлиять на производительность системы.
Фрагменты кода
Рассмотрим несколько фрагментов кода и ожидаемые выходные данные. Обработка исключений игнорируется в этом разделе.
Instantiation
Создание экземпляра надежной параллельной очереди аналогично любой другой надежной коллекции.
IReliableConcurrentQueue<int> queue = await this.StateManager.GetOrAddAsync<IReliableConcurrentQueue<int>>("myQueue");
EnqueueAsync (Асинхронная постановка в очередь)
Ниже приведены несколько фрагментов кода для использования EnqueueAsync, за которым следует их ожидаемые выходные данные.
- Случай 1. Одна задача добавлена в очередь
using (var txn = this.StateManager.CreateTransaction())
{
await this.Queue.EnqueueAsync(txn, 10, cancellationToken);
await this.Queue.EnqueueAsync(txn, 20, cancellationToken);
await txn.CommitAsync();
}
Предположим, что задача выполнена успешно и что параллельные транзакции, изменяющие очередь, не выполнялись. Пользователь может ожидать, что очередь будет содержать элементы в любом из следующих заказов:
10, 20
20, 10
- Задача 2: Параллельная постановка в очередь
// Parallel Task 1
using (var txn = this.StateManager.CreateTransaction())
{
await this.Queue.EnqueueAsync(txn, 10, cancellationToken);
await this.Queue.EnqueueAsync(txn, 20, cancellationToken);
await txn.CommitAsync();
}
// Parallel Task 2
using (var txn = this.StateManager.CreateTransaction())
{
await this.Queue.EnqueueAsync(txn, 30, cancellationToken);
await this.Queue.EnqueueAsync(txn, 40, cancellationToken);
await txn.CommitAsync();
}
Предположим, что задачи успешно выполнены, что задачи выполнялись параллельно и что другие параллельные транзакции не изменяли очередь. Вывод о порядке элементов в очереди не может быть сделан. Для этого фрагмента кода элементы могут отображаться в любом из 4! возможные упорядочения. Очередь пытается сохранить элементы в исходном (вставленном) порядке, но может быть вынуждена их переупорядочить из-за параллельных операций или сбоев.
DequeueAsync
Ниже приведены несколько фрагментов кода для использования TryDequeueAsync, за которым следует ожидаемые выходные данные. Рассмотрим случай, когда очередь уже заполнена следующими элементами.
10, 20, 30, 40, 50, 60
- Случай 1: Одна задача извлечения
using (var txn = this.StateManager.CreateTransaction())
{
await this.Queue.TryDequeueAsync(txn, cancellationToken);
await this.Queue.TryDequeueAsync(txn, cancellationToken);
await this.Queue.TryDequeueAsync(txn, cancellationToken);
await txn.CommitAsync();
}
Предположим, что задача выполнена успешно и что параллельные транзакции, изменяющие очередь, не выполнялись. Так как вывод не может быть сделан по порядку элементов в очереди, любые три элемента могут быть отложены в любом порядке. Очередь пытается сохранить элементы в исходном (поставленном в очередь) порядке, но может быть вынуждена переупорядочивать их из-за одновременных операций или сбоев.
- Случай 2. Параллельная операция dequeue
// Parallel Task 1
List<int> dequeue1;
using (var txn = this.StateManager.CreateTransaction())
{
dequeue1.Add(await this.Queue.TryDequeueAsync(txn, cancellationToken)).val;
dequeue1.Add(await this.Queue.TryDequeueAsync(txn, cancellationToken)).val;
await txn.CommitAsync();
}
// Parallel Task 2
List<int> dequeue2;
using (var txn = this.StateManager.CreateTransaction())
{
dequeue2.Add(await this.Queue.TryDequeueAsync(txn, cancellationToken)).val;
dequeue2.Add(await this.Queue.TryDequeueAsync(txn, cancellationToken)).val;
await txn.CommitAsync();
}
Предположим, что задачи успешно выполнены, что задачи выполнялись параллельно и что другие параллельные транзакции не изменяли очередь. Так как вывод не может быть сделан по порядку элементов в очереди, списки dequeue1 и dequeue2 будут содержать все два элемента в любом порядке.
Один и тот же элемент не будет отображаться в обоих списках. Поэтому, если dequeue1 имеет 10, 30, то dequeue2 будет иметь 20, 40.
- Случай 3. Отмена порядка транзакций
Прерывание транзакции с выполнением удаления из очереди помещает элементы обратно на начало очереди. Порядок, в котором элементы помещаются обратно в начало очереди, не гарантируется. Давайте рассмотрим следующий код:
using (var txn = this.StateManager.CreateTransaction())
{
await this.Queue.TryDequeueAsync(txn, cancellationToken);
await this.Queue.TryDequeueAsync(txn, cancellationToken);
// Abort the transaction
await txn.AbortAsync();
}
Предположим, что элементы были извлечены из очереди в таком порядке:
10, 20
При прерывании транзакции элементы будут добавлены обратно в начало очереди в любом из следующих порядков:
10, 20
20, 10
То же самое верно для всех случаев, когда транзакция не была успешно зафиксирована.
Шаблоны программирования
В этом разделе мы рассмотрим несколько шаблонов программирования, которые могут оказаться полезными при использовании ReliableConcurrentQueue.
Пакетное извлечение из очереди
Рекомендуемый шаблон программирования заключается в том, чтобы задача потребителя выполняла пакетную обработку удалений из очереди, вместо того чтобы выполнять их по одному. Пользователь может выбрать регулирование задержек между каждым пакетом или размером пакета. В следующем фрагменте кода показана эта модель программирования. Помните, что в этом примере обработка выполняется после фиксации транзакции, поэтому если во время обработки произошла ошибка, необработанные элементы будут потеряны без обработки. Кроме того, обработку можно выполнить в пределах области транзакции, однако она может негативно повлиять на производительность и требует обработки уже обработанных элементов.
int batchSize = 5;
long delayMs = 100;
while(!cancellationToken.IsCancellationRequested)
{
// Buffer for dequeued items
List<int> processItems = new List<int>();
using (var txn = this.StateManager.CreateTransaction())
{
ConditionalValue<int> ret;
for(int i = 0; i < batchSize; ++i)
{
ret = await this.Queue.TryDequeueAsync(txn, cancellationToken);
if (ret.HasValue)
{
// If an item was dequeued, add to the buffer for processing
processItems.Add(ret.Value);
}
else
{
// else break the for loop
break;
}
}
await txn.CommitAsync();
}
// Process the dequeues
for (int i = 0; i < processItems.Count; ++i)
{
Console.WriteLine("Value : " + processItems[i]);
}
int delayFactor = batchSize - processItems.Count;
await Task.Delay(TimeSpan.FromMilliseconds(delayMs * delayFactor), cancellationToken);
}
Обработка на основе уведомлений с лучшими усилиями
Другой интересный шаблон программирования использует API Count. Здесь мы можем реализовать уведомительную обработку для очереди с максимальными усилиями. Счётчик очереди можно использовать для регулирования задачи постановки или извлечения из очереди. Обратите внимание, что, как и в предыдущем примере, так как обработка происходит вне транзакции, необработанные элементы могут быть потеряны при возникновении сбоя во время обработки.
int threshold = 5;
long delayMs = 1000;
while(!cancellationToken.IsCancellationRequested)
{
while (this.Queue.Count < threshold)
{
cancellationToken.ThrowIfCancellationRequested();
// If the queue does not have the threshold number of items, delay the task and check again
await Task.Delay(TimeSpan.FromMilliseconds(delayMs), cancellationToken);
}
// If there are approximately threshold number of items, try and process the queue
// Buffer for dequeued items
List<int> processItems = new List<int>();
using (var txn = this.StateManager.CreateTransaction())
{
ConditionalValue<int> ret;
do
{
ret = await this.Queue.TryDequeueAsync(txn, cancellationToken);
if (ret.HasValue)
{
// If an item was dequeued, add to the buffer for processing
processItems.Add(ret.Value);
}
} while (processItems.Count < threshold && ret.HasValue);
await txn.CommitAsync();
}
// Process the dequeues
for (int i = 0; i < processItems.Count; ++i)
{
Console.WriteLine("Value : " + processItems[i]);
}
}
Наилучшее усилие на отвод
Опустошение очереди не гарантируется из-за конкурентной природы структуры данных. Возможно, даже если в очереди нет выполняющихся пользовательских операций, определенный вызов TryDequeueAsync может не вернуть элемент, который ранее был добавлен и подтвержден. Вложенный элемент гарантированно станет видимым для извлечения, однако без внеполосного механизма связи автономный потребитель не может знать, что очередь достигла устойчивого состояния, даже если все производители были остановлены, и новые операции добавления не разрешены. Таким образом, операция сброса выполняется по принципу наилучших усилий, как показано ниже.
Пользователь должен остановить все дальнейшие задачи производителя и потребителя и дождаться фиксации или прерывания любых выполняемых транзакций перед попыткой очистки очереди. Если пользователь знает ожидаемое количество элементов в очереди, он может настроить уведомление, которое сигнализирует о том, что все элементы были удалены.
int numItemsDequeued;
int batchSize = 5;
ConditionalValue ret;
do
{
List<int> processItems = new List<int>();
using (var txn = this.StateManager.CreateTransaction())
{
do
{
ret = await this.Queue.TryDequeueAsync(txn, cancellationToken);
if(ret.HasValue)
{
// Buffer the dequeues
processItems.Add(ret.Value);
}
} while (ret.HasValue && processItems.Count < batchSize);
await txn.CommitAsync();
}
// Process the dequeues
for (int i = 0; i < processItems.Count; ++i)
{
Console.WriteLine("Value : " + processItems[i]);
}
} while (ret.HasValue);
Peek
ReliableConcurrentQueue не предоставляет API TryPeekAsync . Пользователи могут получить семантику просмотра с помощью tryDequeueAsync , а затем прервать транзакцию. В этом примере деквеи обрабатываются только в том случае, если значение элемента больше 10.
using (var txn = this.StateManager.CreateTransaction())
{
ConditionalValue ret = await this.Queue.TryDequeueAsync(txn, cancellationToken);
bool valueProcessed = false;
if (ret.HasValue)
{
if (ret.Value > 10)
{
// Process the item
Console.WriteLine("Value : " + ret.Value);
valueProcessed = true;
}
}
if (valueProcessed)
{
await txn.CommitAsync();
}
else
{
await txn.AbortAsync();
}
}
Необходимо прочитать
- Краткое руководство по надежным службам
- Работа с надежными коллекциями
- Уведомления Reliable Services
- Резервное копирование и восстановление надежных служб (аварийное восстановление)
- Конфигурация Reliable State Manager
- Начало работы со службами веб-API Service Fabric
- Расширенное использование модели программирования надежных служб
- Справочник разработчика по надежным коллекциям