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ą.
GO TO FULL VERSION