1. «Producer‑Consumer» + Concurrent
Chúng ta đã chạm tới pattern «Producer‑Consumer» khi thảo luận về ConcurrentQueue. Hãy xem chi tiết hơn cho trường hợp có nhiều producer và nhiều consumer, cùng với tín hiệu kết thúc chính xác.
Ưu điểm chính của ConcurrentQueue (và các bộ sưu tập Concurrent khác) là nó tự lo việc an toàn luồng. Không cần bọc Enqueue hay TryDequeue trong lock — nhiều luồng có thể tương tác an toàn qua hàng đợi chung.
Ví dụ: Nhiều producer và nhiều consumer
Nhiều worker thread sinh ra nhiệm vụ, và vài thread khác xử lý chúng.
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<string> taskQueue = new ConcurrentQueue<string>();
CancellationTokenSource cts = new CancellationTokenSource(); // Để hủy công việc của các consumer
// Phương thức cho producer
void Producer(string name, int count)
{
for (int i = 0; i < count; i++)
{
string task = $"Zadacha_{name}_{i}";
taskQueue.Enqueue(task);
Console.WriteLine($"[P:{name}] Dobavil: {task}");
Thread.Sleep(10);
}
}
// Phương thức cho consumer
void Consumer(string name)
{
while (!cts.Token.IsCancellationRequested || taskQueue.Count > 0)
{
if (taskQueue.TryDequeue(out string task))
{
Console.WriteLine($"[C:{name}] Obrabotal: {task}");
Thread.Sleep(20);
}
else
{
Thread.Sleep(50); // Đợi nếu hàng đợi rỗng
}
}
Console.WriteLine($"[C:{name}] Zavershil rabotu.");
}
// Chạy ví dụ trong Main:
Task.Run(() => Producer("A", 10)); // Producer A
Task.Run(() => Producer("B", 10)); // Producer B
Task.Run(() => Consumer("1")); // Consumer 1
Task.Run(() => Consumer("2")); // Consumer 2
Thread.Sleep(1000); // Cho thời gian để làm việc
cts.Cancel(); // Báo cho consumer kết thúc công việc
Thread.Sleep(500); // Cho thời gian để consumer lấy phần còn lại và kết thúc
Ở đây vài producer và consumer cùng làm việc với một ConcurrentQueue mà không có race condition: các phương thức Enqueue và TryDequeue là nguyên tử.
Tầm quan trọng của tín hiệu kết thúc công việc (CancellationTokenSource)
Chúng ta dùng CancellationTokenSource (cts) để báo cho các consumer biết khi nào cần kết thúc công việc. Điều này rất quan trọng cho pattern Producer‑Consumer:
- Các producer đã kết thúc. Khi việc thêm phần tử xong, consumer không nên chờ mãi một hàng đợi rỗng vô hạn.
- Ứng dụng kết thúc. Cần dừng consumer một cách chính xác.
CancellationTokenSource và CancellationToken cung cấp cơ chế tiêu chuẩn: consumer kiểm tra định kỳ IsCancellationRequested và khi cần gọi ThrowIfCancellationRequested().
2. BlockingCollection<T>
Mặc dù ConcurrentQueue<T> phù hợp cho “producer‑consumer”, nó yêu cầu chờ thủ công khi hàng đợi rỗng và tự quản lý tín hiệu kết thúc. Để triển khai tiện hơn trong .NET có BlockingCollection<T> — đây không phải là một bộ sưu tập độc lập mà là một wrapper trên bất kỳ IProducerConsumerCollection<T> (ví dụ trên ConcurrentQueue).
Ưu điểm của BlockingCollection:
- Các thao tác blocking. Add()/Take() sẽ block thread nếu collection đầy/rỗng. Không cần tự kiểm tra IsEmpty.
- Giới hạn kích thước. Có thể đặt capacity. Add() sẽ bị block khi đạt giới hạn — tiện để kiểm soát bộ nhớ.
- Kết thúc thuận tiện. CompleteAdding() báo hiệu kết thúc việc thêm, và GetConsumingEnumerable() cho phép consumer xử lý phần tử cho tới khi hoàn toàn xong.
Ví dụ: Producer‑Consumer với BlockingCollection<T>
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
// BlockingCollection theo mặc định dùng ConcurrentQueue
BlockingCollection<int> numbers = new BlockingCollection<int>(capacity: 10); // Hàng đợi giới hạn 10
void ProducerBC(int count)
{
for (int i = 0; i < count; i++)
{
numbers.Add(i); // Block nếu hàng đợi đầy
Console.WriteLine($"[P] Dobavil: {i}");
Thread.Sleep(50);
}
numbers.CompleteAdding(); // Báo rằng producer đã xong
Console.WriteLine("[P] Producer zavershil dobavlenie.");
}
void ConsumerBC()
{
// GetConsumingEnumerable block miễn là có phần tử hoặc chưa gọi CompleteAdding
foreach (var item in numbers.GetConsumingEnumerable())
{
Console.WriteLine($"[C] Obrabotal: {item}");
Thread.Sleep(100);
}
Console.WriteLine("[C] Potrebitel zavershil rabotu.");
}
// Chạy ví dụ trong Main:
Task producerTask = Task.Run(() => ProducerBC(15)); // 15 phần tử, giới hạn 10
Task consumerTask = Task.Run(ConsumerBC);
Task.WaitAll(producerTask, consumerTask); // Chờ kết thúc
Chú ý code của consumer sạch hơn nhờ GetConsumingEnumerable(). Nếu cần thao tác blocking hoặc giới hạn kích thước — BlockingCollection là công cụ cho bạn.
3. Các phương thức và thuộc tính bổ sung của các bộ sưu tập Concurrent
IsEmpty, Count
- IsEmpty (bool): collection có rỗng không.
- Count (int): số phần tử hiện tại.
Ví dụ: Sử dụng IsEmpty và Count
using System.Collections.Concurrent;
ConcurrentQueue<string> q = new ConcurrentQueue<string>();
Console.WriteLine($"Ochered pustaya? {q.IsEmpty}"); // True
q.Enqueue("A");
q.Enqueue("B");
Console.WriteLine($"Elementov v ocheredi: {q.Count}"); // 2
Console.WriteLine($"Ochered pustaya? {q.IsEmpty}"); // False
q.TryDequeue(out var itemA);
Console.WriteLine($"Elementov v ocheredi: {q.Count}"); // 1
Chuyển thành mảng (ToArray())
Tất cả các bộ sưu tập Concurrent cung cấp phương thức ToArray(), trả về một snapshot tức thời của các phần tử.
Ví dụ: Sử dụng ToArray()
using System.Collections.Concurrent;
ConcurrentStack<int> s = new ConcurrentStack<int>();
s.Push(10);
s.Push(20);
s.Push(30);
int[] items = s.ToArray(); // Tạo mảng mới: [30, 20, 10] (vì stack LIFO)
Console.WriteLine($"Elementy v masive: {string.Join(", ", items)}");
// Collection vẫn không thay đổi
Console.WriteLine($"Elementov v stack posle ToArray: {s.Count}"); // 3
Clear collection
Trong .NET 6+ nhiều bộ sưu tập Concurrent có phương thức Clear() để xóa mọi phần tử.
Ví dụ: Xóa collection
using System.Collections.Concurrent;
ConcurrentBag<string> bag = new ConcurrentBag<string>();
bag.Add("Alpha");
bag.Add("Beta");
Console.WriteLine($"Elementov v Bag: {bag.Count}"); // 2
bag.Clear(); // Xóa collection
Console.WriteLine($"Elementov v Bag posle ochistki: {bag.Count}"); // 0
Console.WriteLine($"Bag pust? {bag.IsEmpty}"); // True
4. Những đặc điểm hành vi của các bộ sưu tập Concurrent
Cần nhớ về một snapshot tức thời của dữ liệu. Các thuộc tính như Count và kết quả của ToArray() phản ánh trạng thái của collection tại một thời điểm cụ thể. Trong điều kiện thay đổi song song, giá trị đó có thể lỗi thời ngay sau khi đọc.
Ví dụ: Count và ToArray() — "snapshot"
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);
}
}
// Chạy ví dụ trong Main:
Task.Run(AddItemsContinuously); // Thread liên tục thêm phần tử
Thread.Sleep(100); // Cho chút thời gian thêm
Console.WriteLine($"Tekushchiy Count: {snapshotQueue.Count}"); // Có thể là 50, 80, 120...
Thread.Sleep(100);
Console.WriteLine($"Tekushchiy Count snova: {snapshotQueue.Count}"); // Sẽ khác
int[] currentItems = snapshotQueue.ToArray();
Console.WriteLine($"Kolichestvo elementov v ToArray(): {currentItems.Length}"); // Có thể khác so với Count cuối cùng
Đừng dựa vào Count như một đảm bảo nghiêm ngặt về số phần tử hiện thời khi có thay đổi tích cực.
5. Những lưu ý khi lặp qua các bộ sưu tập Concurrent
Các thao tác riêng lẻ (Add, TryTake, TryPop, GetOrAdd v.v.) là an toàn luồng. Nhưng việc lặp bằng foreach trên một collection đang bị các luồng khác sửa đổi song song không đảm bảo bạn sẽ thấy tất cả phần tử hay chỉ một số — có thể bỏ sót hoặc có hành vi bất ngờ.
Ví dụ: Lặp trong khi sửa đổi
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<int> iterQueue = new ConcurrentQueue<int>();
// Thêm phần tử ban đầu
for (int i = 0; i < 10; i++) iterQueue.Enqueue(i);
// Thread sửa đổi
void Modifier()
{
for (int i = 10; i < 20; i++)
{
iterQueue.Enqueue(i); // Thêm phần tử mới
Thread.Sleep(50);
}
}
// Thread lặp
void Iterator()
{
Console.WriteLine("Nachinayem iteraciyu...");
int count = 0;
foreach (var item in iterQueue) // Lặp
{
Console.Write($"{item} ");
count++;
Thread.Sleep(30); // Giả lập công việc, cho modifier thay đổi collection
}
Console.WriteLine($"\nIteraciya zavershena. Prochitano {count} elementov.");
Console.WriteLine($"Tekuschee kolichestvo v ocheredi: {iterQueue.Count}");
}
// Chạy ví dụ trong Main:
Task.Run(Modifier);
Task.Run(Iterator);
Thread.Sleep(1500); // Cho thời gian chạy
Quy tắc: nếu bạn cần một tập phần tử cố định (ví dụ để báo cáo), trước tiên hãy lấy "snapshot" bằng ToArray(), rồi lặp trên nó:
// Cách đúng để lặp nếu collection có thể thay đổi
int[] snapshot = iterQueue.ToArray();
foreach (var item in snapshot)
{
// Bây giờ bạn lặp trên một mảng-snapshot bất biến
}
Đây là kết thúc phần đi sâu vào các pattern nâng cao và đặc điểm của các bộ sưu tập Concurrent: chúng ta đã xem Producer‑Consumer với nhiều bên tham gia và BlockingCollection, và phân tích các chi tiết quan trọng khi làm việc với Count, ToArray() và việc lặp.
GO TO FULL VERSION