Introdução ao ReliableConcurrentQueue no Azure Service Fabric

A fila simultânea confiável é uma fila assíncrona, transacional e replicada que apresenta alta simultaneidade para operações de enfileiramento e desfila. Foi concebida para oferecer alto débito e baixa latência ao relaxar a ordenação rigorosa de FIFO proporcionada pela Reliable Queue e, em vez disso, oferece uma ordenação de melhor esforço.

APIs

Fila simultânea Fila simultânea confiável
void Enqueue(item T) Task EnqueueAsync(ITransaction tx, item T)
bool TentarDesenfileirar(out T resultado) Tarefa< ConditionalValue < T >> TryDequeueAsync(ITransaction tx)
int Contagem() contagem longa()

Comparação com fila confiável

A Fila Simultânea Confiável é oferecida como uma alternativa à Fila Confiável. Deve ser usada em casos em que não é exigida uma ordenação rigorosa de FIFO, pois garantir o FIFO exige um comprometimento com a concorrência. Reliable Queue utiliza bloqueios para impor a ordem FIFO, com no máximo uma transação permitida para enfileirar e no máximo uma transação permitida para desenfileirar de cada vez. Em comparação, a Fila Confiável Simultânea relaxa a restrição de ordenação e permite a qualquer número de transações simultâneas intercalar as suas operações de enfileirar e desenfileirar. A ordem de melhor esforço é fornecida, no entanto, a ordenação relativa de dois valores em uma fila simultânea confiável nunca pode ser garantida.

A Fila Simultânea Confiável fornece maior taxa de transferência e menor latência do que a Fila Confiável sempre que há várias transações simultâneas executando enfileiramentos e/ou desfilas.

Um exemplo de caso de uso para o ReliableConcurrentQueue é o cenário Message Queue . Nesse cenário, um ou mais produtores de mensagens criam e adicionam itens à fila e um ou mais consumidores de mensagens extraem mensagens da fila e as processam. Vários produtores e consumidores podem trabalhar de forma independente, usando transações simultâneas para processar a fila.

Diretrizes de uso

  • A fila espera que os itens na fila tenham um período de retenção baixo. Ou seja, os itens não ficariam na fila durante muito tempo.
  • A fila não garante uma ordenação estrita FIFO.
  • A fila não lê as suas próprias escritas. Se um item estiver enfileirado dentro de uma transação, ele não será visível para um desfilador dentro da mesma transação.
  • As dequeues não estão isoladas entre si. Se o item A for retirado da fila na transação txnA, mesmo que txnA não esteja confirmado, o item A não será visível para uma transação concorrente txnB. Se txnA abortar, A se tornará visível para txnB imediatamente.
  • O comportamento TryPeekAsync pode ser implementado usando um TryDequeueAsync e, em seguida, anulando a transação. Um exemplo desse comportamento pode ser encontrado na seção Padrões de programação.
  • A contagem não é transacional. Pode ser usado para obter uma ideia do número de elementos na fila, mas representa um ponto no tempo e não pode ser fiável.
  • O processamento dispendioso dos itens retirados da fila não deve ser realizado enquanto a transação está ativa, para evitar transações de longa duração que possam ter impacto no desempenho do sistema.

Trechos de código

Vejamos alguns trechos de código e suas saídas esperadas. O tratamento de exceções é ignorado nesta seção.

Instanciação

A criação de uma instância de uma Fila Confiável Concorrente é semelhante à de qualquer outra Coleção Confiável.

IReliableConcurrentQueue<int> queue = await this.StateManager.GetOrAddAsync<IReliableConcurrentQueue<int>>("myQueue");

EnqueueAsync

Aqui estão alguns trechos de código para usar EnqueueAsync seguido por suas saídas esperadas.

  • Caso 1: Tarefa de fila única
using (var txn = this.StateManager.CreateTransaction())
{
    await this.Queue.EnqueueAsync(txn, 10, cancellationToken);
    await this.Queue.EnqueueAsync(txn, 20, cancellationToken);

    await txn.CommitAsync();
}

Suponha que a tarefa foi concluída com êxito e que não houve transações simultâneas modificando a fila. O usuário pode esperar que a fila contenha os itens em qualquer uma das seguintes ordens:

10, 20

20, 10

  • Caso 2: Tarefa de Enfileiramento Paralelo
// 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();
}

Suponha que as tarefas foram concluídas com êxito, que as tarefas foram executadas em paralelo e que não houve outras transações simultâneas modificando a fila. Nenhuma inferência pode ser feita sobre a ordem dos itens na fila. Para este trecho de código, os itens podem aparecer em qualquer um dos 4! possíveis encomendas. A fila tenta manter os itens na ordem original (enfileirada), mas pode ser forçada a reordená-los devido a operações ou falhas simultâneas.

DequeueAsync

Aqui estão alguns trechos de código para usar TryDequeueAsync seguido pelas saídas esperadas. Suponha que a fila já está preenchida com os seguintes itens na fila:

10, 20, 30, 40, 50, 60

  • Caso 1: Tarefa única de desfila
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();
}

Suponha que a tarefa foi concluída com êxito e que não houve transações simultâneas modificando a fila. Como nenhuma inferência pode ser feita sobre a ordem dos itens na fila, qualquer um dos três itens pode ser retirado da fila, em qualquer ordem. A fila tenta manter os itens na ordem original (enfileirada), mas pode ser forçada a reordená-los devido a operações ou falhas simultâneas.

  • Caso 2: Tarefa de desfila paralela
// 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();
}

Suponha que as tarefas foram concluídas com êxito, que as tarefas foram executadas em paralelo e que não houve outras transações simultâneas modificando a fila. Como nenhuma inferência pode ser feita sobre a ordem dos itens na fila, as listas dequeue1 e dequeue2 conterão, cada uma, dois itens, em qualquer ordem.

O mesmo item não* aparece em ambas as listas. Assim, se dequeue1 tem 10, 30, então dequeue2 teria 20, 40.

  • Caso 3: Ordem de Retirada com Aborto de Transação

Abortar uma transação com desenfileiramentos em andamento recoloca os itens no início da fila. A ordem em que os itens são colocados de volta no topo da fila não é garantida. Vejamos o seguinte código:

using (var txn = this.StateManager.CreateTransaction())
{
    await this.Queue.TryDequeueAsync(txn, cancellationToken);
    await this.Queue.TryDequeueAsync(txn, cancellationToken);

    // Abort the transaction
    await txn.AbortAsync();
}

Suponha que os itens foram retirados da fila na seguinte ordem:

10, 20

Quando abortamos a transação, os itens são adicionados de volta ao topo da fila em qualquer uma das seguintes ordens:

10, 20

20, 10

O mesmo se aplica a todos os casos em que a transação não foi Confirmada com sucesso.

Padrões de programação

Nesta seção, vamos examinar alguns padrões de programação que podem ser úteis no uso de ReliableConcurrentQueue.

Desfilas em lote

Um padrão de programação recomendado é que a tarefa do consumidor agrupe as suas retiradas de fila em vez de executar uma retirada de cada vez. O usuário pode optar por limitar os atrasos entre cada lote ou o tamanho do lote. O trecho de código a seguir mostra esse modelo de programação. Esteja ciente, neste exemplo, o processamento é feito depois que a transação é confirmada, portanto, se ocorrer uma falha durante o processamento, os itens não processados serão perdidos sem terem sido processados. Alternativamente, o processamento pode ser feito dentro do escopo da transação, no entanto, pode ter um impacto negativo no desempenho e requer o manuseio dos itens já processados.

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);
}

Processamento Best-Effort Notification-Based

Outro padrão de programação interessante usa a API Count. Aqui, podemos implementar o processamento baseado em notificações de esforço máximo para a fila. A Contagem da fila pode ser usada para limitar uma operação de enqueue ou dequeue. Observe que, como no exemplo anterior, como o processamento ocorre fora da transação, os itens não processados podem ser perdidos se ocorrer uma falha durante o processamento.

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]);
    }
}

Best-Effort Dreno

Não se pode garantir um esvaziamento da fila devido à natureza concorrente da estrutura de dados. É possível que, mesmo que nenhuma operação de utilizador na fila esteja em andamento, uma determinada chamada para TryDequeueAsync possa não devolver um item que já foi enfileirado e confirmado. O item introduzido na fila está garantido a eventualmente tornar-se visível para a remoção da fila; no entanto, sem um mecanismo de comunicação externo, um consumidor independente não pode saber que a fila atingiu um estado estável, mesmo que todos os produtores tenham sido parados e não são permitidas novas operações de inclusão na fila. Assim, a operação de drenagem é considerada de melhor esforço, conforme implementado abaixo.

O usuário deve parar todas as outras tarefas de produtor e consumidor e aguardar que qualquer transação em voo seja confirmada ou abortada, antes de tentar drenar a fila. Se o usuário souber o número esperado de itens na fila, ele poderá configurar uma notificação que sinalize que todos os itens foram retirados da fila.

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);

Espreitar

O ReliableConcurrentQueue não fornece a API do TryPeekAsync . Os usuários podem obter a semântica de visualização usando um TryDequeueAsync e, em seguida, abortando a transação. Neste exemplo, os itens são removidos da fila e processados apenas se o valor do item for maior que 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();
    }
}

Leitura obrigatória