You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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.Run in C#

Here's a practical, thread-safe implementation that meets all your needs:

1. Define the Cache Structure

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());
    }
}
2. Thread-Safe Cache Manager

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));
    }
}
3. Worker Thread Implementation

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
        }
    }
}
4. Usage Example

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.");
    }
}
Key Explanations
  • 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.Exchange ensures cache updates are thread-safe without blocking locks.
  • Non-Blocking Behavior: Task.Run is used to start the worker and handle client updates asynchronously, so no thread gets blocked waiting for operations to complete.
  • Graceful Shutdown: The CancellationToken lets us stop the worker cleanly, and it's okay if client updates don't finish before exit—no data corruption risk.

内容的提问来源于stack exchange,提问作者allmhuran

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 04:24:44