open_files_op.cc 8.3 KB
Newer Older
F
fengjiayi 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14
//   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.

15
#include <cmath>
Y
yuyang18 已提交
16
#include <stdexcept>
F
fengjiayi 已提交
17
#include <thread>  // NOLINT
18 19
#include "ThreadPool.h"
#include "paddle/fluid/framework/blocking_queue.h"
20
#include "paddle/fluid/operators/reader/blocking_queue.h"
Y
yuyang18 已提交
21
#include "paddle/fluid/operators/reader/buffered_reader.h"
F
fengjiayi 已提交
22 23 24 25 26 27
#include "paddle/fluid/operators/reader/reader_op_registry.h"

namespace paddle {
namespace operators {
namespace reader {

28
class IReaderContainer {
F
fengjiayi 已提交
29
 public:
30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46
  virtual ~IReaderContainer() {}
  virtual void AppendReader(
      std::unique_ptr<framework::ReaderBase>&& readers) = 0;
  virtual void Stop() = 0;
  virtual void Start() = 0;
  virtual void ReadNext(std::vector<framework::LoDTensor>* out) = 0;
};

class OrderedReaderContainer : public IReaderContainer {
 public:
  void AppendReader(std::unique_ptr<framework::ReaderBase>&& reader) override {
    pending_.emplace(std::move(reader));
  }

  void Stop() override {
    while (!pending_.empty()) {
      MoveFrontPendingToDone();
47
    }
F
fengjiayi 已提交
48 49
  }

50
  void Start() override { std::swap(done_, pending_); }
F
fengjiayi 已提交
51

52 53 54 55 56 57 58 59 60 61 62
  void ReadNext(std::vector<framework::LoDTensor>* out) override {
    if (!pending_.empty()) {
      pending_.front()->ReadNext(out);
      if (out->empty()) {
        MoveFrontPendingToDone();
        ReadNext(out);
      }
    } else {
      out->clear();
    }
  }
F
fengjiayi 已提交
63

F
fengjiayi 已提交
64
 private:
65 66 67 68 69 70 71 72 73
  void MoveFrontPendingToDone() {
    pending_.front()->Shutdown();
    pending_.front()->Start();
    done_.emplace(move(pending_.front()));
    pending_.pop();
  }

  std::queue<std::unique_ptr<framework::ReaderBase>> pending_;
  std::queue<std::unique_ptr<framework::ReaderBase>> done_;
F
fengjiayi 已提交
74 75
};

76 77
class PreemptiveReaderContainer : public IReaderContainer {
  using ReaderList = std::list<std::unique_ptr<framework::ReaderBase>>;
F
fengjiayi 已提交
78

79 80 81
  struct FutureItem {
    std::vector<framework::LoDTensor> data_;
    ReaderList::iterator reader_it_;
Y
yuyang18 已提交
82
    std::exception_ptr exception_;
83
  };
F
fengjiayi 已提交
84

85
  using FutureList = std::list<std::future<FutureItem>>;
F
fengjiayi 已提交
86

87 88
 public:
  explicit PreemptiveReaderContainer(size_t thread_num) : pool_(thread_num) {}
F
fengjiayi 已提交
89

90 91 92 93 94 95 96 97 98 99 100 101
  void Stop() override {
    if (!pending_.empty()) {
      for (auto& reader : pending_) {
        reader->Shutdown();
      }
      for (auto& fu : futures_) {
        fu.wait();
      }
      futures_.clear();
      for (auto& reader : pending_) {
        reader->Start();
        done_.emplace_back(std::move(reader));
F
fengjiayi 已提交
102
      }
103 104 105 106
      pending_.clear();
      bool timeout;
      complete_queue_.PopAll(1000, &timeout);
      PADDLE_ENFORCE(!timeout);
F
fengjiayi 已提交
107 108
    }
  }
109 110 111 112

  void Start() override {
    for (auto& reader : done_) {
      AppendReader(std::move(reader));
F
fengjiayi 已提交
113
    }
114
    done_.clear();
F
fengjiayi 已提交
115
  }
116 117 118 119 120

  void ReadNext(std::vector<framework::LoDTensor>* out) override {
    if (!pending_.empty()) {
      auto future_it = complete_queue_.Pop();
      FutureItem item = future_it->get();
Y
yuyang18 已提交
121 122 123 124 125 126 127 128 129
      if (item.exception_) {
        for (auto it = futures_.begin(); it != futures_.end(); ++it) {
          if (it != future_it) {
            it->wait();  // Wait all other threads complete.
          }
        }
        std::rethrow_exception(item.exception_);

      } else if (item.data_.empty()) {  // reader done.
130 131 132 133 134 135 136 137 138 139 140
        done_.emplace_back(std::move(*item.reader_it_));
        pending_.erase(item.reader_it_);
        futures_.erase(future_it);
        ReadNext(out);
      } else {
        *out = item.data_;
        // continue read async
        AsyncRead(item.reader_it_, &future_it);
      }
    } else {
      out->clear();
F
fengjiayi 已提交
141
    }
142 143 144
  }

 private:
Y
yuyang18 已提交
145 146
  void AppendReader(std::unique_ptr<framework::ReaderBase>&& reader) override {
    pending_.emplace_back(std::move(reader));
147 148 149 150 151 152 153 154 155 156 157 158 159 160
    auto reader_it = pending_.end();
    --reader_it;

    futures_.emplace_back();
    auto future_it = futures_.end();
    --future_it;

    AsyncRead(reader_it, &future_it);
  }

  void AsyncRead(const ReaderList::iterator& reader_it,
                 FutureList::iterator* future_it_ptr) {
    auto& future_it = *future_it_ptr;
    *future_it = pool_.enqueue([reader_it, future_it, this] {
Y
yuyang18 已提交
161 162 163 164 165 166 167 168 169 170 171 172 173 174 175
      try {
        FutureItem item;
        item.reader_it_ = reader_it;
        (*reader_it)->ReadNext(&item.data_);
        if (item.data_.empty()) {
          (*reader_it)->Shutdown();
          (*reader_it)->Start();
        }
        complete_queue_.Push(future_it);
        return item;
      } catch (...) {
        FutureItem item;
        item.exception_ = std::current_exception();
        complete_queue_.Push(future_it);
        return item;
176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193
      }
    });
  }

  FutureList futures_;
  ThreadPool pool_;
  framework::BlockingQueue<FutureList::iterator> complete_queue_;
  std::list<std::unique_ptr<framework::ReaderBase>> pending_;
  std::list<std::unique_ptr<framework::ReaderBase>> done_;
};

class MultiFileReader : public framework::ReaderBase {
 public:
  MultiFileReader(const std::vector<std::string>& file_names,
                  std::unique_ptr<IReaderContainer>&& container)
      : container_(std::move(container)) {
    for (auto& fn : file_names) {
      container_->AppendReader(CreateReaderByFileName(fn));
F
fengjiayi 已提交
194 195
    }
  }
196

197 198 199 200 201
  ~MultiFileReader() { container_->Stop(); }

 protected:
  void ReadNextImpl(std::vector<framework::LoDTensor>* out) override {
    container_->ReadNext(out);
F
fengjiayi 已提交
202
  }
203 204 205 206 207 208
  void ShutdownImpl() override { container_->Stop(); }
  void StartImpl() override { container_->Start(); }

 private:
  std::unique_ptr<IReaderContainer> container_;
};
F
fengjiayi 已提交
209 210 211 212 213 214 215 216 217 218 219 220

class OpenFilesOp : public framework::OperatorBase {
 public:
  using framework::OperatorBase::OperatorBase;

 private:
  void RunImpl(const framework::Scope& scope,
               const platform::Place& dev_place) const override {
    const auto& shape_concat = Attr<std::vector<int>>("shape_concat");
    const auto& ranks = Attr<std::vector<int>>("ranks");
    PADDLE_ENFORCE(!shape_concat.empty() && !ranks.empty());
    PADDLE_ENFORCE_EQ(std::accumulate(ranks.begin(), ranks.end(), 0),
F
fengjiayi 已提交
221
                      static_cast<int>(shape_concat.size()),
F
fengjiayi 已提交
222 223 224 225
                      "The accumulate of all ranks should be equal to the "
                      "shape concat's length.");
    const auto& file_names = Attr<std::vector<std::string>>("file_names");
    PADDLE_ENFORCE(!file_names.empty(), "No file to be read!");
226
    bool is_test = Attr<bool>("is_test");
F
fengjiayi 已提交
227 228 229

    auto* out = scope.FindVar(Output("Out"))
                    ->template GetMutable<framework::ReaderHolder>();
230 231 232 233 234 235
    std::unique_ptr<IReaderContainer> container;

    if (is_test) {
      container.reset(new OrderedReaderContainer());
    } else {
      container.reset(new PreemptiveReaderContainer(
Y
yuyang18 已提交
236
          static_cast<size_t>(Attr<int>("thread_num"))));
237 238
    }

Y
yuyang18 已提交
239 240 241 242 243 244 245 246
    std::shared_ptr<framework::ReaderBase> reader(
        new MultiFileReader(file_names, std::move(container)));
    auto buffer_size = Attr<int>("buffer_size");
    if (buffer_size > 1) {
      reader = framework::MakeDecoratedReader<BufferedReader>(
          reader, platform::CPUPlace(), buffer_size);
    }
    out->Reset(reader);
F
fengjiayi 已提交
247 248 249
  }
};

250
class OpenFilesOpMaker : public FileReaderMakerBase {
Y
Yu Yang 已提交
251 252
 protected:
  void Apply() override {
253
    AddAttr<std::vector<std::string>>("file_names", "Files to be read.");
254
    AddAttr<bool>("is_test", "Used for testing data.").SetDefault(false);
255

F
fengjiayi 已提交
256 257 258
    AddComment(R"DOC(
      OpenFiles Operator

Y
Yu Yang 已提交
259
      An OpenFilesOp creates a MultiFileReader, which is able to
F
fengjiayi 已提交
260 261
      read data multi-threaded from multiple files.
    )DOC");
Y
yuyang18 已提交
262 263 264 265 266
    AddAttr<int>("thread_num",
                 "The maximal concurrent prefetch thread number. Used only "
                 "when is_test = False");
    AddAttr<int>("buffer_size", "The reading buffer of these files.")
        .GreaterThan(0);
F
fengjiayi 已提交
267 268 269 270 271 272 273 274 275 276
  }
};

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

namespace reader = paddle::operators::reader;

REGISTER_FILE_READER_OPERATOR(open_files, reader::OpenFilesOp,
277
                              reader::OpenFilesOpMaker);