【问题标题】:flushing thread local buffer at end of parallel loop with TBB使用 TBB 在并行循环结束时刷新线程本地缓冲区
【发布时间】:2016-12-14 19:30:57
【问题描述】:

我想并行化一个循环(使用tbb),其中包含一些昂贵但可矢量化的迭代(随机分布)。我的想法是缓冲它们并在达到矢量大小时刷新缓冲区。这样的缓冲区必须是线程本地的。例如,

// dummy for testing
void do_vectorized_work(size_t k, size_t*indices)
{}
// dummy for testing
bool requires_expensive_work(size_t k)
{ return (k&7)==0; }

struct buffer
{
  size_t K=0, B[vector_size];
  void load(size_t i)
  {
    B[K++]=i;
    if(K==vector_size)
      flush();
  }
  void flush()
  {
    do_vectorized_work(K,B);
    K=0;
  }
};

void do_work_in_parallel(size_t N)
{
  tbb::enumerable_thread_specific<buffer> tl_buffer;

  tbb::parallel_for(size_t(0),N,[&](size_t i)
  {
    if(requires_expensive_work(i))
      tl_buffer.local().load(i);
  });
}

但是,这会使缓冲区不为空,所以我仍然需要最后一次刷新每个缓冲区

for(auto&b:tl_buffer)
  b.flush();

但这是连续剧!当然,我也可以尝试并行这样做

using tl_range = typename tbb::enumerable_thread_specific<buffer>::range_type;
tbb::parallel_for(tl_buffer.range(),[](tl_range const&range)
{
  for(auto r:range)
    r->flush();
});

但我不确定这是否有效(因为缓冲区的数量与线程的数量一样多)。我想知道是否有可能在事件发生后避免最后的冲洗。 IE。是否可以使用tbb::tasks(替换tbb::parallel_for)以使每个线程的最终任务是刷新其缓冲区?

【问题讨论】:

    标签: c++ multithreading parallel-processing tbb


    【解决方案1】:

    不,工作线程没有关于这个特定任务是否是给定工作的最后一个任务的完整信息(这就是工作窃取的工作方式)。因此,不可能在parallel_for 或调度程序本身的级别上实现这样的功能。因此,我建议您使用您描述的这两种方法。

    不过,您还可以做两件事。

    • 使其异步。 IE。将一项任务排入队列,该任务将刷新所有内容。这将有助于从主线程的热路径中删除此代码。如果在完成此任务时需要设置任何依赖项,请小心。
    • 使用tbb::task_scheduler_observer 来初始化线程特定的数据,并在线程关闭或一段时间内没有工作时延迟释放它。后者需要使用 local observer feature,它尚未得到官方支持,但已经稳定了几年。

    例子:

    #define TBB_PREVIEW_LOCAL_OBSERVER 1
    #include <tbb/tbb.h>
    #include <assert.h>
    
    typedef void * buffer_t;
    const static int bufsz = 1024;
    class thread_buffer_allocator: public tbb::task_scheduler_observer {
      tbb::enumerable_thread_specific<buffer_t> _buf;
    public:
      thread_buffer_allocator( )
        : tbb::task_scheduler_observer( /*local=*/ true ) {
        observe(true); // activate the observer
      }
      ~thread_buffer_allocator( ) {
        observe(false); // deactivate the observer
        for(auto &b : _buf) {
            printf("destructor: cleared: %p\n", b);
            free(b);
        }
      }
      /*override*/ void on_scheduler_entry( bool worker ) {
        assert(_buf.local() == nullptr);
        _buf.local() = malloc(bufsz);
        printf("on entry: %p\n", _buf.local());
      }
      /*override*/ void on_scheduler_exit( bool worker ) {
        printf("on exit\n");
        if(_buf.local()) {
            printf("on exit: cleared %p\n", _buf.local());
            free(_buf.local());
            _buf.local() = nullptr;
        }
      }
    };
    
    int main() {
      thread_buffer_allocator buffers_scope;
      tbb::parallel_for(0, 1024*1024*1024, [&](auto i){
        usleep(i%3);
      });
      return 0;
    }
    

    【讨论】:

    • 谢谢。我不认为异步方法比我在 OP 中描述的尝试更好。使用tbb::task_scheduler_observer 的方法听起来很有趣。你能用代码 sn-p 概述一下这是如何工作的吗?
    • @Walter 已更新。虽然我只在本地观察者没有足够近期 TBB 的在线编译器上尝试过:coliru.stacked-crooked.com/a/11728cd935579cfe
    【解决方案2】:

    我突然想到,这可以通过减少来解决。

    struct buffer
    {
      std::size_t K=0, B[vector_size];
      void load(std::size_t i)
      {
        B[K++]=i;
        if(K==vector_size) flush();
      }
      void flush()
      {
        do_vectorized_work(K,B);
        K=0;
      }
      buffer(buffer const&, tbb::split)
      {}
      void operator()(tbb::block_range<std::size_t> const&range)
      { for(i:range) load(i); }
      bool empty()
      { return K==0; }
      std::size_t pop()
      { return K? B[--K] : 0; }
      void join(buffer&rhs)
      { while(!rhs.empty()) load(rhs.pop()); }
    };
    
    void do_work_in_parallel(std::size_t N)
    {
      buffer buff;
      tbb::parallel_reduce(tbb::block_range<std::size_t>(0,N,vector_size),buff);
      if(!buff.empty())
        buff.flush();
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-08-10
      • 2012-04-02
      • 1970-01-01
      • 2021-02-04
      • 1970-01-01
      相关资源
      最近更新 更多