buffered_reader.cc 5.5 KB
Newer Older
Y
yuyang18 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
// Copyright (c) 2018 PaddlePaddle Authors. All Rights Reserved.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
//     http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

#include "paddle/fluid/operators/reader/buffered_reader.h"
C
chengduo 已提交
16
#include <memory>
Y
yuyang18 已提交
17
#include <vector>
D
Dun Liang 已提交
18
#include "paddle/fluid/framework/data_type.h"
Y
yuyang18 已提交
19

20
#include "paddle/fluid/platform/profiler.h"
Y
yuyang18 已提交
21 22 23
namespace paddle {
namespace operators {
namespace reader {
F
fengjiayi 已提交
24 25 26 27 28 29
BufferedReader::~BufferedReader() {
  reader_->Shutdown();
  while (!position_.empty()) {
    position_.front().wait();
    position_.pop();
  }
D
Dun Liang 已提交
30 31 32 33
#ifdef PADDLE_WITH_CUDA
  if (platform::is_gpu_place(place_)) {
    platform::SetDeviceId(boost::get<platform::CUDAPlace>(place_).device);
    PADDLE_ENFORCE(cudaStreamDestroy(stream));
D
Dun Liang 已提交
34
    for (auto &event : events) PADDLE_ENFORCE(cudaEventDestroy(event));
D
Dun Liang 已提交
35 36
  }
#endif
F
fengjiayi 已提交
37 38
}

Y
yuyang18 已提交
39 40 41 42 43 44 45
BufferedReader::BufferedReader(
    const std::shared_ptr<framework::ReaderBase> &reader,
    const platform::Place &place, size_t buffer_size)
    : framework::DecoratedReader(reader),
      thread_pool_(1),
      place_(place),
      buffer_size_(buffer_size) {
D
Dun Liang 已提交
46 47 48
#ifdef PADDLE_WITH_CUDA
  if (platform::is_gpu_place(place_)) {
    platform::SetDeviceId(boost::get<platform::CUDAPlace>(place_).device);
D
Dun Liang 已提交
49 50 51 52 53
    compute_stream =
        ((platform::CUDADeviceContext *)(platform::DeviceContextPool::Instance()
                                             .Get(place_)))
            ->stream();
    events.resize(buffer_size);
54
    for (auto &event : events) {
D
Dun Liang 已提交
55
      PADDLE_ENFORCE(cudaEventCreateWithFlags(&event, cudaEventDisableTiming));
56
    }
57
    PADDLE_ENFORCE(cudaStreamCreateWithFlags(&stream, cudaStreamNonBlocking));
D
Dun Liang 已提交
58 59
  }
#endif
Y
yuyang18 已提交
60 61
  cpu_buffer_.resize(buffer_size);
  gpu_buffer_.resize(buffer_size);
Y
yuyang18 已提交
62
  ReadTillBufferFullAsync();
Y
yuyang18 已提交
63
}
F
fengjiayi 已提交
64

Y
yuyang18 已提交
65
void BufferedReader::ReadTillBufferFullAsync() {
Y
yuyang18 已提交
66 67
  PADDLE_ENFORCE_EQ(position_.size(), 0U);
  for (size_t i = 0; i < buffer_size_; ++i) {
Y
yuyang18 已提交
68
    ReadAsync(i);
Y
yuyang18 已提交
69 70
  }
}
F
fengjiayi 已提交
71

Y
yuyang18 已提交
72
void BufferedReader::ReadAsync(size_t i) {
D
Dun Liang 已提交
73 74 75 76 77 78
#ifdef PADDLE_WITH_CUDA
  if (platform::is_gpu_place(place_)) {
    platform::SetDeviceId(boost::get<platform::CUDAPlace>(place_).device);
    PADDLE_ENFORCE(cudaEventRecord(events[i], compute_stream));
  }
#endif
Y
yuyang18 已提交
79 80 81
  position_.emplace(thread_pool_.enqueue([this, i]() -> size_t {
    TensorVec &cpu = cpu_buffer_[i];
    reader_->ReadNext(&cpu);
Y
yuyang18 已提交
82

Y
yuyang18 已提交
83 84 85
    if (cpu.empty()) {
      return -1UL;
    }
Y
yuyang18 已提交
86

D
Dun Liang 已提交
87 88
#ifdef PADDLE_WITH_CUDA
    // NOTE(liangdun): using async copy instead of TensorCopySync
89 90 91
    // TensorCopySync would block other stream, because TensorCopySync
    // issues the copying command to the default stream, it will make two
    // commands from different streams cannot run concurrently.
Y
yuyang18 已提交
92
    if (platform::is_gpu_place(place_)) {
D
Dun Liang 已提交
93 94
      platform::SetDeviceId(boost::get<platform::CUDAPlace>(place_).device);
      PADDLE_ENFORCE(cudaStreamWaitEvent(stream, events[i], 0));
Y
yuyang18 已提交
95 96
      TensorVec &gpu = gpu_buffer_[i];
      gpu.resize(cpu.size());
97
      platform::RecordEvent record_event("BufferedReader:MemoryCopy");
Y
yuyang18 已提交
98
      for (size_t i = 0; i < cpu.size(); ++i) {
D
Dun Liang 已提交
99 100 101 102 103 104 105
        gpu[i].Resize(cpu[i].dims());
        gpu[i].set_layout(cpu[i].layout());
        auto cpu_place = cpu[i].place();
        auto cpu_ptr = cpu[i].data<void>();
        auto gpu_ptr = gpu[i].mutable_data(place_, cpu[i].type());
        auto size =
            cpu[i].numel() * paddle::framework::SizeOfType(cpu[i].type());
106
        if (platform::is_cuda_pinned_place(cpu_place)) {
D
Dun Liang 已提交
107 108 109
          memory::Copy(boost::get<platform::CUDAPlace>(place_), gpu_ptr,
                       boost::get<platform::CUDAPinnedPlace>(cpu_place),
                       cpu_ptr, size, stream);
110
        } else if ((platform::is_gpu_place(cpu_place))) {
D
Dun Liang 已提交
111 112 113
          memory::Copy(boost::get<platform::CUDAPlace>(place_), gpu_ptr,
                       boost::get<platform::CUDAPlace>(cpu_place), cpu_ptr,
                       size, stream);
114 115
        } else {
          // TODO(zcd): The default stream should not be used here.
D
Dun Liang 已提交
116 117
          memory::Copy(boost::get<platform::CUDAPlace>(place_), gpu_ptr,
                       boost::get<platform::CPUPlace>(cpu_place), cpu_ptr, size,
118
                       0);
119
        }
Y
yuyang18 已提交
120
        gpu[i].set_lod(cpu[i].lod());
Y
yuyang18 已提交
121
      }
D
Dun Liang 已提交
122
      PADDLE_ENFORCE(cudaStreamSynchronize(stream));
Y
yuyang18 已提交
123
    }
D
Dun Liang 已提交
124
#endif
Y
yuyang18 已提交
125
    return i;
Y
yuyang18 已提交
126 127
  }));
}
F
fengjiayi 已提交
128

Y
yuyang18 已提交
129 130
void BufferedReader::ShutdownImpl() {
  reader_->Shutdown();
Y
yuyang18 已提交
131 132 133
  while (!position_.empty()) {
    position_.pop();
  }
Y
yuyang18 已提交
134
  prev_pos_ = -1UL;
Y
yuyang18 已提交
135
}
F
fengjiayi 已提交
136

Y
yuyang18 已提交
137 138
void BufferedReader::StartImpl() {
  reader_->Start();
Y
yuyang18 已提交
139
  ReadTillBufferFullAsync();
Y
yuyang18 已提交
140
}
F
fengjiayi 已提交
141

Y
yuyang18 已提交
142
void BufferedReader::ReadNextImpl(std::vector<framework::LoDTensor> *out) {
Y
yuyang18 已提交
143 144 145 146 147 148 149 150 151 152 153 154 155
  if (position_.empty()) {
    out->clear();
    return;
  }
  size_t i = position_.front().get();
  position_.pop();

  if (i == -1UL) {
    ReadNextImpl(out);
    return;
  }

  *out = platform::is_gpu_place(place_) ? gpu_buffer_[i] : cpu_buffer_[i];
Y
yuyang18 已提交
156 157 158 159 160 161 162 163

  // Do not push current position into ReadAsync. Push the previous position
  // Since all computation in fluid are async, change the data of
  // current position may cause data error.
  if (prev_pos_ != -1Ul) {
    ReadAsync(prev_pos_);
  }
  prev_pos_ = i;
Y
yuyang18 已提交
164 165 166 167 168
}

}  // namespace reader
}  // namespace operators
}  // namespace paddle