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

寻求结合ConcurrentQueue与ConcurrentStack的双边访问并发集合方案

Great question! You're looking for a concurrent collection that blends the best of ConcurrentQueue (FIFO for cleaning old items) and ConcurrentStack (LIFO-style access to the newest items) — let's break down how to build this step by step.

Core Requirements Recap

First, let's restate your needs to make sure we're aligned:

  • Multi-producer safe: Multiple threads can add elements without conflicts.
  • Access latest n elements: Clients should retrieve the most recently generated items, ordered from newest to oldest.
  • Auto-cleanup old items: Remove any elements older than x hours, starting from the oldest (FIFO-style cleanup).
Implementation Approach

The key challenge is supporting efficient, thread-safe operations on both ends: adding new items to the "newest" end, cleaning old items from the "oldest" end, and fetching the latest n items.

We'll use ConcurrentLinkedQueue as the underlying storage (it's thread-safe for enqueue/dequeue operations) and wrap it with logic to handle time-based cleanup and recent item retrieval. Here's a complete, production-ready implementation in C#:

Code Example: TimedConcurrentBuffer<T>

using System;
using System.Collections.Generic;
using System.Collections.Concurrent;
using System.Threading;

public class TimedConcurrentBuffer<T> : IDisposable
{
    private readonly ConcurrentLinkedQueue<(T Item, DateTime Timestamp)> _itemQueue = new();
    private readonly TimeSpan _retentionWindow;
    private readonly Timer _cleanupTimer;
    private readonly object _cleanupSyncLock = new();

    // Initialize with retention period and cleanup interval (e.g., every 5 minutes)
    public TimedConcurrentBuffer(TimeSpan retentionWindow, TimeSpan cleanupInterval)
    {
        _retentionWindow = retentionWindow;
        _cleanupTimer = new Timer(TriggerCleanup, null, TimeSpan.Zero, cleanupInterval);
    }

    // Thread-safe add operation for multi-producer scenarios
    public void AddItem(T item)
    {
        _itemQueue.Enqueue((Item: item, Timestamp: DateTime.UtcNow));
    }

    // Retrieve the latest N elements, ordered from newest to oldest
    public IEnumerable<T> GetLatestItems(int count)
    {
        // Ensure we're not returning stale items by running a quick cleanup first
        TriggerCleanup();

        var recentItems = new List<T>();
        foreach (var queueItem in _itemQueue)
        {
            recentItems.Add(queueItem.Item);
            // Stop collecting once we have enough items
            if (recentItems.Count == count) break;
        }

        // Queue stores oldest -> newest, so reverse to get newest -> oldest order
        recentItems.Reverse();
        return recentItems;
    }

    // Clean up items older than the retention window
    private void TriggerCleanup()
    {
        // Use a lock to prevent concurrent cleanup attempts (avoids redundant work)
        if (!Monitor.TryEnter(_cleanupSyncLock))
        {
            return; // Another thread is already cleaning up, skip
        }

        try
        {
            var cutoffTime = DateTime.UtcNow.Subtract(_retentionWindow);
            // Keep removing old items from the queue head until we hit a valid one
            while (_itemQueue.TryPeek(out var oldestItem) && oldestItem.Timestamp < cutoffTime)
            {
                _itemQueue.TryDequeue(out _);
            }
        }
        finally
        {
            Monitor.Exit(_cleanupSyncLock);
        }
    }

    // Clean up resources when done
    public void Dispose()
    {
        _cleanupTimer?.Dispose();
    }
}
How It Works

Let's walk through the key components:

1. Multi-Producer Safety

  • ConcurrentLinkedQueue handles thread-safe enqueues out of the box — multiple threads can call AddItem simultaneously without race conditions.

2. Retrieving Latest N Items

  • The queue stores items from oldest to newest (head = oldest, tail = newest).
  • We traverse the queue to collect items, then reverse the list to return newest-to-oldest order.
  • If there are fewer than count items in the buffer, we return all available items.

3. Auto-Cleanup of Old Items

  • A Timer runs periodic cleanup (configurable interval) to remove items older than your retention window.
  • Cleanup happens from the queue head (oldest items first), which matches the FIFO-style cleanup you need.
  • A lock ensures only one thread runs cleanup at a time, avoiding redundant checks and potential race conditions.
Optimization Tips
  • Adjust Cleanup Interval: For high-throughput systems, use a longer cleanup interval (e.g., 10 minutes) to reduce lock contention. For low-volume systems, you could even trigger cleanup on every AddItem or GetLatestItems call instead of using a timer.
  • Strong Consistency (Optional): If you need guaranteed consistency between cleanup and retrieval, wrap the GetLatestItems logic in the same lock used for cleanup. This will reduce concurrent performance but eliminate any chance of returning stale items during a cleanup.
  • Large Retention Windows: If your buffer grows very large, consider adding a maximum size limit (e.g., keep only the last 10,000 items) alongside the time-based retention to prevent memory bloat.
Usage Example
// Create a buffer that keeps items for 2 hours, cleaning up every 5 minutes
var buffer = new TimedConcurrentBuffer<string>(TimeSpan.FromHours(2), TimeSpan.FromMinutes(5));

// Multi-producer scenario: multiple threads adding items
new Thread(() => buffer.AddItem($"Log entry 1: {DateTime.Now}")).Start();
new Thread(() => buffer.AddItem($"Log entry 2: {DateTime.Now}")).Start();

// Client retrieves the latest 5 items
var latestLogs = buffer.GetLatestItems(5);
foreach (var log in latestLogs)
{
    Console.WriteLine(log);
}

// Don't forget to dispose when done
buffer.Dispose();

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:09:06