compute_interceptor.h 1.8 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
// Copyright (c) 2021 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.

#pragma once

17
#include <queue>
18 19
#include <utility>

20 21 22 23 24
#include "paddle/fluid/distributed/fleet_executor/interceptor.h"

namespace paddle {
namespace distributed {

25 26
const int64_t INFINITE_BUFFER_SIZE = -1;

27 28 29 30
class ComputeInterceptor : public Interceptor {
 public:
  ComputeInterceptor(int64_t interceptor_id, TaskNode* node);

31 32 33 34
 protected:
  virtual void RunOps();
  virtual void SendDataReadyToDownStream();
  virtual void ReplyCompletedToUpStream();
35 36 37 38 39 40
  virtual void Compute(const InterceptorMessage& msg);
  void Run();
  void IncreaseReady(int64_t up_id, int64_t scope_id);
  void DecreaseBuff(int64_t down_id);

  int64_t cur_scope_id_;
41

42 43 44 45 46
  // upstream_id-->(max_ready_size, scope-->ready_size)
  std::map<int64_t, std::pair<int64_t, std::map<int64_t, int64_t>>>
      in_readys_{};
  // downstream_id-->(max_buffer_size, used_size)
  std::map<int64_t, std::pair<int64_t, int64_t>> out_buffs_{};
47 48

 private:
49
  void PrepareDeps();
50 51
  InterceptorMessage PrepareVarsMsg();
  void DecodeMsgVars(const InterceptorMessage& msg);
52

53 54
  bool IsInputReady();
  bool CanWriteOutput();
55
  std::map<int64_t, bool> scope_id_to_finish_flag_;
56 57 58 59
};

}  // namespace distributed
}  // namespace paddle