【问题标题】:Pattern for concurrent cache sharing并发缓存共享模式
【发布时间】:2011-05-22 12:48:03
【问题描述】:

好的,我有点不确定如何最好地命名这个问题 :) 但是假设这种情况,你是 出去并获取一些网页(带有各种网址)并将其缓存在本地。即使使用多个线程,缓存部分也很容易解决。

但是,假设一个线程开始获取一个 url,几毫秒后另一个线程想要获取相同的 url。有没有什么好的模式可以让秒线程的方法在第一个线程上等待以获取页面,将其插入缓存并返回它,这样您就不必执行多个请求。即使对于大约需要 300-700 毫秒的请求,开销也足够少?并且没有锁定对其他 url 的请求

基本上,当对相同 url 的请求彼此紧接时,我希望第二个请求“搭载”第一个请求

当您开始获取页面并锁定它时,我有一些松散的想法,即有一个字典,您可以在其中插入一个带有键作为 url 的对象。如果有任何匹配的键已经是对象,则锁定它,然后尝试获取实际缓存的 url。

我有点不确定细节,但是为了使其真正线程安全,使用 ConcurrentDictionary 可能是其中的一部分...

这样的场景有什么通用的模式和解决方案吗?

分解错误行为:

线程 1:检查缓存,它不存在所以开始获取 url

线程 2:开始获取相同的 url,因为它仍然不存在于缓存中

线程1:完成并插入缓存,返回页面

线程 2:完成并插入缓存(或丢弃),返回页面

分解正确行为:

线程 1:检查缓存,它不存在所以开始获取 url

线程 2:想要相同的 url,但看到它当前正在被获取,所以等待线程 1

线程1:完成并插入缓存,返回页面

线程 2:通知线程 1 已完成并返回它获取的页面线程 1

编辑

到目前为止,大多数解决方案似乎都误解了这个问题,只解决了缓存问题,正如我所说的那样,这不是问题,问题是在进行外部 Web 抓取以进行第二次抓取时在第一次抓取之前完成缓存它以使用第一个结果而不是第二个

【问题讨论】:

  • 我的回答确实解决了您在编辑中提出的问题。
  • @Luke,您当前的解决方案似乎确实是我正在寻找的,谢谢!我会等几个小时等待任何替代解决方案,然后我会关闭问题
  • 您是否考虑过一种解决方案,您可以使用某种同步字典(例如 ConcurrentDictionary),其中 url 作为键,IAsyncResult 作为值?如果线程 2 将尝试获取当前正在由线程 1 下载的页面,则它只需要等待 IAsyncResult 直到它完成然后获取页面内容(IAsyncResult 可能不是正确的选择,但你会得到想法...)。

标签: c# multithreading design-patterns concurrency c#-4.0


【解决方案1】:

您可以使用ConcurrentDictionary<K,V> 和double-checked locking 的变体:

public static string GetUrlContent(string url)
{
    object value1 = _cache.GetOrAdd(url, new object());

    if (value1 == null)    // null check only required if content
        return null;       // could legitimately be a null string

    var urlContent = value1 as string;
    if (urlContent != null)
        return urlContent;    // got the content

    // value1 isn't a string which means that it's an object to lock against
    lock (value1)
    {
        object value2 = _cache[url];

        // at this point value2 will *either* be the url content
        // *or* the object that we already hold a lock against
        if (value2 != value1)
            return (string)value2;    // got the content

        urlContent = FetchContentFromTheWeb(url);    // todo
        _cache[url] = urlContent;
        return urlContent;
    }
}

private static readonly ConcurrentDictionary<string, object> _cache =
                                  new ConcurrentDictionary<string, object>();

【讨论】:

    【解决方案2】:

    编辑:我的代码现在有点难看,但每个 URL 使用单独的锁。这允许异步获取不同的 URL,但是每个 URL 只会被获取一次。

    public class UrlFetcher
    {
        static Hashtable cache = Hashtable.Synchronized(new Hashtable());
    
        public static String GetCachedUrl(String url)
        {
            // exactly 1 fetcher is created per URL
            InternalFetcher fetcher = (InternalFetcher)cache[url];
            if( fetcher == null )
            {
                lock( cache.SyncRoot )
                {
                    fetcher = (InternalFetcher)cache[url];
                    if( fetcher == null )
                    {
                        fetcher = new InternalFetcher(url);
                        cache[url] = fetcher;
                    }
                }
            }
            // blocks all threads requesting the same URL
            return fetcher.Contents;
        }
    
        /// <summary>Each fetcher locks on itself and is initilized with null contents.
        /// The first thread to call fetcher.Contents will cause the fetch to occur, and
        /// block until completion.</summary>
        private class InternalFetcher
        {
            private String url;
            private String contents;
    
            public InternalFetcher(String url)
            {
                this.url = url;
                this.contents = null;
            }
    
            public String Contents
            {
                get
                {
                    if( contents == null )
                    {
                        lock( this ) // "this" is an instance of InternalFetcher...
                        {
                            if( contents == null )
                            {
                                contents = FetchFromWeb(url);
                            }
                        }
                    }
                    return contents;
                }
            }
        }
    }
    

    【讨论】:

    • 我也使用该模式进行实际缓存。但是,如果对 url 的请求紧随其后,并且还没有被缓存,它并不能解决问题
    • 第二个请求直到第一个释放它才获得锁,所以它会看到缓存中的项目。
    • 查看更新后的描述,您正在为所有请求锁定缓存,与其说是缓存不如说是获取
    • 好的,我想我明白你在说什么。运行 2 个并发 fetch 请求是可以的,只要它们针对不同的 URL。我会更新我的代码:)
    • 好吧,但是如果有有个对同一个url的并发请求,第二个请求应该等待并返回第一个请求
    【解决方案3】:

    Semaphore请站起来!起来!站起来!

    使用Semaphore,您可以轻松地与它同步您的线程。 在这两种情况下

    1. 您正在尝试加载当前正在缓存的页面
    2. 您正在将缓存保存到从中加载页面的文件中。

    在这两种情况下你都会遇到麻烦。

    就像写者和读者问题一样,是操作系统竞速问题中的常见问题。仅当线程想要重建缓存或开始缓存页面时,任何线程都不应从中读取。如果一个线程正在读取它,它应该等到读取完成并替换缓存,没有2个线程应该将相同的页面缓存到同一个文件中。因此,所有读者都可以随时从缓存中读取,因为没有写者在上面写。

    你应该在 msdn 上阅读一些使用示例的信号量,它非常易于使用。只是想要做某事的线程调用信号量,如果资源可以授予它执行工作,否则会休眠并等待资源准备好时被唤醒。

    【讨论】:

    • 啊,也许我有点不清楚,我需要根据请求的 url 缓存/共享...我可以使用基于键的信号量吗?
    【解决方案4】:

    免责声明:这可能是一个愚蠢的答案。请原谅我,如果是的话。

    我建议使用一些带锁的共享字典对象来跟踪当前获取或已经获取的 url。

    • 在每次请求时,对照该对象检查 url。

    • 如果存在 url 条目,请检查缓存。 (这意味着其中一个线程已经获取或正在获取它)

    • 如果它在缓存中可用,则使用它,否则让当前线程休眠一段时间并再次检查。 (如果不在缓存中,则某些线程仍在获取它,所以等待它完成)

    • 如果在字典对象中找不到该条目,则将 url 添加到它并发送请求。获得响应后,将其添加到缓存中。

    这个逻辑应该可以工作,但是,您需要注意缓存过期和从字典对象中删除条目。

    【讨论】:

      【解决方案5】:

      我的解决方案是在缓存超时或不存在时使用 atomicBoolean 控制访问数据库;

      同时,只有一个线程(我称之为read-th)可以访问数据库,其他线程自旋直到read-th返回数据并将其写入缓存;

      这里是代码; java实现;

      public class CacheBreakDownDefender<K, R> {
      
      /**
       * false = do not write null to cache when get null value from database;
       */
      private final boolean writeNullToCache;
      
      /**
       * cache different query key
       */
      private final ConcurrentHashMap<K, AtomicBoolean> selectingDBTagMap = new ConcurrentHashMap<>();
      
      
      public static <K, R> CacheBreakDownDefender<K, R> getInstance(Class<K> keyType, Class<R> resultType) {
          return Singleton.get(keyType.getName() + resultType.getName(), () -> new CacheBreakDownDefender<>(false));
      }
      
      public static <K, R> CacheBreakDownDefender<K, R> getInstance(Class<K> keyType, Class<R> resultType, boolean writeNullToCache) {
          return Singleton.get(keyType.getName() + resultType.getName(), () -> new CacheBreakDownDefender<>(writeNullToCache));
      }
      
      private CacheBreakDownDefender(boolean writeNullToCache) {
          this.writeNullToCache = writeNullToCache;
      }
      
      public R readFromCache(K key, Function<K, ? extends R> getFromCache, Function<K, ? extends R> getFromDB, BiConsumer<K, R> writeCache) throws InterruptedException {
          R result = getFromCache.apply(key);
          if (result == null) {
              final AtomicBoolean selectingDB = selectingDBTagMap.computeIfAbsent(key, x -> new AtomicBoolean(false));
              if (selectingDB.compareAndSet(false, true)) { 
                  try { 
                      result = getFromDB.apply(key);
                      if (result != null || writeNullToCache) {
                          writeCache.accept(key, result);
                      }
                  } finally {
                      selectingDB.getAndSet(false);
                      selectingDBTagMap.remove(key);
                  }
              } else {
                  
                  while (selectingDB.get()) {
                      TimeUnit.MILLISECONDS.sleep(0L);
                      //do nothing...  
                  }
                  return getFromCache.apply(key);
              }
          }
          return result;
      }
      
      public static void main(String[] args) throws InterruptedException {
      
          Map<String, String> map = new ConcurrentHashMap<>();
          CacheBreakDownDefender<String, String> instance = CacheBreakDownDefender.getInstance(String.class, String.class, true);
      
          for (int i = 0; i < 9; i++) {
              int finalI = i;
              new Thread(() -> {
                  String kele = null;
                  try {
                      if (finalI == 6) {
                          kele = instance.readFromCache("kele2", map::get, key -> "helloword2", map::put);
                      } else
                          kele = instance.readFromCache("kele", map::get, key -> "helloword", map::put);
                  } catch (InterruptedException e) {
                      Thread.currentThread().interrupt();
                  }
                  log.info("resut= {}", kele);
              }).start();
          }
          TimeUnit.SECONDS.sleep(2L);
      }
      

      }

      【讨论】:

        【解决方案6】:

        这并不完全适用于并发缓存,而是适用于所有缓存:

        "A cache with a bad policy is another name for a memory leak" (Raymond Chen)

        【讨论】:

          猜你喜欢
          • 2015-06-25
          • 1970-01-01
          • 2015-11-05
          • 1970-01-01
          • 2011-01-04
          • 2023-03-07
          • 2011-07-27
          • 1970-01-01
          • 2013-04-02
          相关资源
          最近更新 更多