C#中使用Task.Run实现非阻塞锁及多线程共享缓存处理问询
Alright, let's tackle your problem step by step. You've got a shared cache accessed by at least two threads: a worker thread running a continuous loop to use cache data, and client threads that submit new information to the cache. The core requirements are:
- The worker uses the new data only on its next iteration (after finishing the current loop's work)
- No harm if the program exits before all client writes complete
- Implement non-blocking behavior with
Task.Runin C#
Here's a practical, thread-safe implementation that meets all your needs:
First, create an immutable cache wrapper (or use .NET's built-in ImmutableDictionary for better memory efficiency) to ensure we never modify a cache instance once it's created—this eliminates race conditions during worker processing.
public class AppCache { // Use ImmutableDictionary for built-in thread safety and immutability public ImmutableDictionary<string, object> Data { get; } public AppCache(ImmutableDictionary<string, object> data) { Data = data ?? ImmutableDictionary<string, object>.Empty; } // Helper to create a cache from a regular dictionary public static AppCache FromDictionary(Dictionary<string, object> source) { return new AppCache(source.ToImmutableDictionary()); } }
This class handles atomic updates to the cache and provides non-blocking access for both workers and clients. We use Interlocked.Exchange for thread-safe reference swapping—no blocking locks needed here.
public class CacheManager { // Volatile ensures threads see the latest cache reference immediately private volatile AppCache _currentCache = new AppCache(ImmutableDictionary<string, object>.Empty); // Get the latest cache snapshot (non-blocking, zero overhead) public AppCache GetCurrentSnapshot() { return _currentCache; } // Submit new data (atomic, non-blocking) public void UpdateCache(Dictionary<string, object> newData) { var newCache = AppCache.FromDictionary(newData); // Atomically swap the cache reference—no two threads can do this at the same time Interlocked.Exchange(ref _currentCache, newCache); } // Async wrapper using Task.Run for non-blocking client calls public Task UpdateCacheAsync(Dictionary<string, object> newData) { // Offload the (lightweight) cache creation to a background thread // so the calling thread (e.g., UI thread) doesn't block return Task.Run(() => UpdateCache(newData)); } }
The worker runs in a continuous loop, grabbing a cache snapshot at the start of each iteration. This guarantees it finishes processing the current snapshot before picking up any new updates from clients.
public class CacheWorker { private readonly CacheManager _cacheManager; private CancellationTokenSource _cancellationSource = new(); public CacheWorker(CacheManager cacheManager) { _cacheManager = cacheManager; } // Start the worker in a non-blocking way using Task.Run public void Start() { Task.Run(() => RunWorkerLoop(_cancellationSource.Token), _cancellationSource.Token); } // Gracefully stop the worker public void Stop() { _cancellationSource.Cancel(); } private async Task RunWorkerLoop(CancellationToken token) { while (!token.IsCancellationRequested) { // Grab the latest cache snapshot ONLY at the start of the loop var currentCache = _cacheManager.GetCurrentSnapshot(); // Process the cache data—this can be a long-running operation await ProcessCacheData(currentCache, token); // Optional: Add a delay to avoid spiking CPU usage await Task.Delay(1000, token); } } private async Task ProcessCacheData(AppCache cache, CancellationToken token) { // Replace this with your actual business logic Console.WriteLine($"Processing cache with {cache.Data.Count} entries..."); foreach (var (key, value) in cache.Data) { if (token.IsCancellationRequested) break; Console.WriteLine($"- {key}: {value}"); await Task.Delay(200, token); // Simulate work } } }
Putting it all together in a console app to simulate client updates and worker processing:
class Program { static async Task Main(string[] args) { var cacheManager = new CacheManager(); var worker = new CacheWorker(cacheManager); // Start the worker (non-blocking) worker.Start(); Console.WriteLine("Worker started. Press any key to exit..."); // Simulate client threads submitting new data var client1 = Task.Run(async () => { await Task.Delay(2500); // Wait 2.5 seconds var newData = new Dictionary<string, object> { { "User:1", "Alice" }, { "User:2", "Bob" } }; await cacheManager.UpdateCacheAsync(newData); Console.WriteLine("Client 1 submitted new cache data"); }); var client2 = Task.Run(async () => { await Task.Delay(5000); // Wait 5 seconds var newData = new Dictionary<string, object> { { "User:3", "Charlie" }, { "User:1", "Alice Updated" } }; await cacheManager.UpdateCacheAsync(newData); Console.WriteLine("Client 2 submitted new cache data"); }); // Wait for user input to exit Console.ReadKey(); // Clean up worker.Stop(); await Task.WhenAll(client1, client2); Console.WriteLine("Program exited."); } }
- Immutable Cache Instances: By never modifying an existing cache, we ensure the worker never sees partial updates—once it grabs a snapshot, it's guaranteed to stay consistent until the next loop iteration.
- Atomic Reference Swaps:
Interlocked.Exchangeensures cache updates are thread-safe without blocking locks. - Non-Blocking Behavior:
Task.Runis used to start the worker and handle client updates asynchronously, so no thread gets blocked waiting for operations to complete. - Graceful Shutdown: The
CancellationTokenlets us stop the worker cleanly, and it's okay if client updates don't finish before exit—no data corruption risk.
内容的提问来源于stack exchange,提问作者allmhuran

