CodeGym /Kursy /C# SELF /Zaawansowane wzorce i cechy kolekcji

Zaawansowane wzorce i cechy kolekcji Concurrent

C# SELF
Poziom 58 , Lekcja 3
Dostępny

1. „Producent‑Konsument” + Concurrent

Już poruszaliśmy wzorzec „Producent‑Konsument” przy omawianiu ConcurrentQueue. Przyjrzyjmy się mu dokładniej dla przypadku z wieloma producentami i konsumentami oraz poprawnym sygnalizowaniem zakończenia.

Główną zaletą ConcurrentQueue (i innych kolekcji Concurrent) jest to, że sama dba o bezpieczeństwo wątków. Nie trzeba opakowywać Enqueue czy TryDequeue w lock — kilka wątków może bezpiecznie współdziałać przez wspólną kolejkę.

Przykład: Kilku producentów i kilku konsumentów

Kilka wątków roboczych generuje zadania, a inne wątki je przetwarzają.

using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

ConcurrentQueue<string> taskQueue = new ConcurrentQueue<string>();
CancellationTokenSource cts = new CancellationTokenSource(); // Do anulowania pracy konsumentów

// Metoda producenta
void Producer(string name, int count)
{
    for (int i = 0; i < count; i++)
    {
        string task = $"Zadanie_{name}_{i}";
        taskQueue.Enqueue(task);
        Console.WriteLine($"[P:{name}] Dodał: {task}");
        Thread.Sleep(10); 
    }
}

// Metoda konsumenta
void Consumer(string name)
{
    while (!cts.Token.IsCancellationRequested || taskQueue.Count > 0)
    {
        if (taskQueue.TryDequeue(out string task))
        {
            Console.WriteLine($"[C:{name}] Przetworzył: {task}");
            Thread.Sleep(20); 
        }
        else
        {
            Thread.Sleep(50); // Czekamy, jeśli kolejka jest pusta
        }
    }
    Console.WriteLine($"[C:{name}] Zakończył pracę.");
}

// Uruchomienie przykładu w Main:
Task.Run(() => Producer("A", 10)); // Producent A
Task.Run(() => Producer("B", 10)); // Producent B
Task.Run(() => Consumer("1"));    // Konsument 1
Task.Run(() => Consumer("2"));    // Konsument 2

Thread.Sleep(1000); // Dajemy czas na pracę
cts.Cancel();       // Sygnał dla konsumentów, żeby zakończyli pracę
Thread.Sleep(500); // Dajemy czas konsumentom, żeby dokończyli resztki i zakończyli pracę

Tutaj kilku producentów i konsumenci pracują jednocześnie z jedną ConcurrentQueue bez wyścigów: metody Enqueue i TryDequeue są atomowe.

Znaczenie sygnałów zakończenia pracy (CancellationTokenSource)

Używamy CancellationTokenSource (cts) do sygnalizowania konsumentom potrzeby zakończenia pracy. To krytyczne dla wzorca Producer‑Consumer:

  • Producenci zakończyli pracę. Gdy dodawanie elementów się skończy, konsumenci nie powinni czekać w nieskończoność na pustą kolejkę.
  • Aplikacja się zamyka. Trzeba poprawnie zatrzymać konsumentów.

CancellationTokenSource i CancellationToken dają standardowy mechanizm: konsument okresowo sprawdza IsCancellationRequested i w razie potrzeby wywołuje ThrowIfCancellationRequested().

2. BlockingCollection<T>

Choć ConcurrentQueue<T> dobrze nadaje się do „producent‑konsument”, wymaga ręcznego oczekiwania przy pustej kolejce i samodzielnego sygnalizowania zakończenia. Do wygodniejszej implementacji w .NET jest BlockingCollection<T> — to nie jest oddzielna kolekcja, a wrapper nad dowolnym IProducerConsumerCollection<T> (np. nad ConcurrentQueue).

Zalety BlockingCollection:

  • Operacje blokujące. Add()/Take() blokują wątek, jeśli kolekcja jest pełna/pusta. Nie trzeba ręcznie sprawdzać IsEmpty.
  • Ograniczenie rozmiaru. Można ustawić pojemność (Capacity). Add() zablokuje się, gdy limit zostanie osiągnięty — wygodne do kontroli pamięci.
  • Wygodne finalizowanie. CompleteAdding() sygnalizuje zakończenie dodawania, a GetConsumingEnumerable() pozwala konsumentowi przetwarzać elementy aż do całkowitego zakończenia.

Przykład: Producer‑Consumer z BlockingCollection<T>

using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

// BlockingCollection domyślnie używa ConcurrentQueue
BlockingCollection<int> numbers = new BlockingCollection<int>(capacity: 10); // Kolejka z limitem 10

void ProducerBC(int count)
{
    for (int i = 0; i < count; i++)
    {
        numbers.Add(i); // Blokuje, jeśli kolejka jest pełna
        Console.WriteLine($"[P] Dodał: {i}");
        Thread.Sleep(50);
    }
    numbers.CompleteAdding(); // Sygnalizujemy, że producent skończył
    Console.WriteLine("[P] Producent zakończył dodawanie.");
}

void ConsumerBC()
{
    // GetConsumingEnumerable blokuje, dopóki są elementy lub dopóki nie wywołano CompleteAdding
    foreach (var item in numbers.GetConsumingEnumerable())
    {
        Console.WriteLine($"[C] Przetworzył: {item}");
        Thread.Sleep(100);
    }
    Console.WriteLine("[C] Konsument zakończył pracę.");
}

// Uruchomienie przykładu w Main:
Task producerTask = Task.Run(() => ProducerBC(15)); // 15 elementów, limit 10
Task consumerTask = Task.Run(ConsumerBC);
Task.WaitAll(producerTask, consumerTask); // Czekamy na zakończenie

Zwróć uwagę, jak czyściejszy jest kod konsumenta dzięki GetConsumingEnumerable(). Jeśli potrzebujesz operacji blokujących lub ograniczenia rozmiaru — BlockingCollection będzie dobrym narzędziem.

3. Dodatkowe metody i właściwości kolekcji Concurrent

IsEmpty, Count

  • IsEmpty (bool): czy kolekcja jest pusta.
  • Count (int): bieżąca liczba elementów.

Przykład: Użycie IsEmpty i Count

using System.Collections.Concurrent;

ConcurrentQueue<string> q = new ConcurrentQueue<string>();
Console.WriteLine($"Czy kolejka jest pusta? {q.IsEmpty}"); // True

q.Enqueue("A");
q.Enqueue("B");
Console.WriteLine($"Elementów w kolejce: {q.Count}"); // 2
Console.WriteLine($"Czy kolejka jest pusta? {q.IsEmpty}"); // False

q.TryDequeue(out var itemA);
Console.WriteLine($"Elementów w kolejce: {q.Count}"); // 1

Konwersja do tablic (ToArray())

Wszystkie kolekcje Concurrent udostępniają metodę ToArray(), która zwraca migawkę elementów w danym momencie.

Przykład: Użycie ToArray()

using System.Collections.Concurrent;

ConcurrentStack<int> s = new ConcurrentStack<int>();
s.Push(10);
s.Push(20);
s.Push(30);

int[] items = s.ToArray(); // Stworzy nową tablicę: [30, 20, 10] (dla stosu LIFO)
Console.WriteLine($"Elementy w tablicy: {string.Join(", ", items)}");

// Kolekcja pozostaje niezmieniona
Console.WriteLine($"Elementów w stosie po ToArray: {s.Count}"); // 3

Czyszczenie kolekcji

W .NET 6+ wiele kolekcji Concurrent dostało metodę Clear() do usuwania wszystkich elementów.

Przykład: Czyszczenie kolekcji

using System.Collections.Concurrent;

ConcurrentBag<string> bag = new ConcurrentBag<string>();
bag.Add("Alpha");
bag.Add("Beta");
Console.WriteLine($"Elementów w Bag: {bag.Count}"); // 2

bag.Clear(); // Czyścimy kolekcję
Console.WriteLine($"Elementów w Bag po wyczyszczeniu: {bag.Count}"); // 0
Console.WriteLine($"Bag puste? {bag.IsEmpty}"); // True

4. Cechy zachowania kolekcji Concurrent

Ważne, żeby pamiętać o migawce danych. Właściwości takie jak Count i wyniki ToArray() odzwierciedlają stan kolekcji w konkretnym momencie. W warunkach równoległych zmian ta wartość może stać się nieaktualna dosłownie zaraz po odczycie.

Przykład: Count i ToArray() — „migawka”

using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

ConcurrentQueue<int> snapshotQueue = new ConcurrentQueue<int>();

void AddItemsContinuously()
{
    for (int i = 0; i < 1000; i++)
    {
        snapshotQueue.Enqueue(i);
        Thread.Sleep(1); 
    }
}

// Uruchomienie przykładu w Main:
Task.Run(AddItemsContinuously); // Wątek, który ciągle dodaje elementy

Thread.Sleep(100); // Dajemy trochę czasu na dodawanie
Console.WriteLine($"Aktualny Count: {snapshotQueue.Count}"); // Może być 50, 80, 120...
Thread.Sleep(100);
Console.WriteLine($"Aktualny Count ponownie: {snapshotQueue.Count}"); // Będzie inna wartość
int[] currentItems = snapshotQueue.ToArray();
Console.WriteLine($"Liczba elementów w ToArray(): {currentItems.Length}"); // Może różnić się od ostatniego Count

Nie polegaj na Count jako na ścisłej gwarancji aktualnej liczby elementów podczas aktywnych zmian.

5. Niuanse iteracji po kolekcjach Concurrent

Pojedyncze operacje (Add, TryTake, TryPop, GetOrAdd itp.) są bezpieczne dla wątków. Ale iteracja z użyciem foreach po kolekcji, która jest równolegle modyfikowana przez inne wątki, nie gwarantuje, że zobaczysz wszystkie elementy lub wyłącznie je — mogą wystąpić pominięcia i niespodzianki.

Przykład: Iteracja podczas modyfikacji

using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;

ConcurrentQueue<int> iterQueue = new ConcurrentQueue<int>();

// Dodajemy początkowe elementy
for (int i = 0; i < 10; i++) iterQueue.Enqueue(i);

// Wątek-modyfikator
void Modifier()
{
    for (int i = 10; i < 20; i++)
    {
        iterQueue.Enqueue(i); // Dodajemy nowe elementy
        Thread.Sleep(50);
    }
}

// Wątek-iterator
void Iterator()
{
    Console.WriteLine("Zaczynamy iterację...");
    int count = 0;
    foreach (var item in iterQueue) // Iterujemy
    {
        Console.Write($"{item} ");
        count++;
        Thread.Sleep(30); // Symulacja pracy, dajemy modyfikatorowi zmienić kolekcję
    }
    Console.WriteLine($"\nIteracja zakończona. Odczytano {count} elementów.");
    Console.WriteLine($"Aktualna liczba w kolejce: {iterQueue.Count}");
}

// Uruchomienie przykładu w Main:
Task.Run(Modifier);
Task.Run(Iterator);
Thread.Sleep(1500); // Dajemy czas na pracę

Zasada: jeśli potrzebujesz stałego zbioru elementów (np. do raportu), najpierw zrób „migawkę” przez ToArray(), a potem iteruj po niej:

// Poprawny sposób iteracji, jeśli kolekcja może się zmieniać
int[] snapshot = iterQueue.ToArray();
foreach (var item in snapshot)
{
    // Teraz iterujesz po niezmiennym tablicowym snapshotcie
}

To kończy nasze zanurzenie w zaawansowane wzorce i cechy kolekcji Concurrent: obejrzeliśmy Producer‑Consumer z wieloma uczestnikami i BlockingCollection, oraz omówiliśmy ważne niuanse pracy z Count, ToArray() i iteracją.

2
Zadanie
C# SELF, poziom 58, lekcja 3
Niedostępne
Przykład użycia `BlockingCollection` z ograniczeniem rozmiaru
Przykład użycia `BlockingCollection` z ograniczeniem rozmiaru
Komentarze
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION