task_node.cc 3.6 KB
Newer Older
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
// 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.

#include "paddle/fluid/distributed/fleet_executor/task_node.h"
16
#include "paddle/fluid/framework/op_registry.h"
17 18 19 20 21 22 23 24
#include "paddle/fluid/framework/operator.h"

namespace paddle {
namespace distributed {
namespace {
using OperatorBase = TaskNode::OperatorBase;
}

L
LiYuRio 已提交
25 26 27 28 29 30 31 32 33
TaskNode::TaskNode(const framework::ProgramDesc& program, int64_t rank,
                   int64_t max_run_times, int64_t max_slot_nums)
    : program_(program),
      rank_(rank),
      max_run_times_(max_run_times),
      max_slot_nums_(max_slot_nums) {
  // Should be serially invoked, not thread-safe
  static int64_t task_node_cnt = 0;
  task_id_ = task_node_cnt++;
34 35 36 37 38 39
  for (const auto& op_desc : program.Block(0).AllOps()) {
    ops_vec_.emplace_back(framework::OpRegistry::CreateOp(*op_desc));
  }
  for (const auto& op : ops_vec_) {
    ops_.emplace_back(op.get());
  }
L
LiYuRio 已提交
40 41
}

42
TaskNode::TaskNode(int32_t role, const std::vector<OperatorBase*>& ops,
43 44 45 46 47 48 49 50
                   int64_t rank, int64_t task_id, int64_t max_run_times,
                   int64_t max_slot_nums)
    : ops_(ops),
      role_(role),
      rank_(rank),
      task_id_(task_id),
      max_run_times_(max_run_times),
      max_slot_nums_(max_slot_nums) {}
51

52
TaskNode::TaskNode(int32_t role, int64_t rank, int64_t task_id,
53 54 55 56 57 58
                   int64_t max_run_times, int64_t max_slot_nums)
    : role_(role),
      rank_(rank),
      task_id_(task_id),
      max_run_times_(max_run_times),
      max_slot_nums_(max_slot_nums) {}
59

60 61 62
bool TaskNode::AddUpstreamTask(int64_t task_id, int64_t buff_size) {
  const auto& ret = upstream_.emplace(task_id, buff_size);
  return ret.second;
L
LiYuRio 已提交
63
}
64

65 66 67
bool TaskNode::AddDownstreamTask(int64_t task_id, int64_t buff_size) {
  const auto& ret = downstream_.emplace(task_id, buff_size);
  return ret.second;
68
}
69 70 71 72 73 74 75 76 77 78

std::string TaskNode::DebugString() const {
  std::ostringstream os;
  os << "role: " << role_ << ", task_id: " << task_id_ << "\n";
  for (std::size_t i = 0; i < ops_.size(); ++i) {
    os << ops_[i]->Type() << " ";
  }
  os << "\n";
  return os.str();
}
79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107

void TaskNode::SetRunPerSteps(int64_t value) {
  PADDLE_ENFORCE_GE(value, 1,
                    platform::errors::InvalidArgument(
                        "run_per_steps must >= 1, but received %ld", value));
  run_per_steps_ = value;
}

void TaskNode::SetRunAtOffset(int64_t value) {
  PADDLE_ENFORCE_GE(value, 0,
                    platform::errors::InvalidArgument(
                        "run_at_offset must >= 0, but received %ld", value));
  run_at_offset_ = value;
}

void TaskNode::SetReplyUpPerSteps(int64_t value) {
  PADDLE_ENFORCE_GE(
      value, 1, platform::errors::InvalidArgument(
                    "reply_up_per_steps must >= 1, but received %ld", value));
  reply_up_per_steps_ = value;
}

void TaskNode::SetSendDownPerSteps(int64_t value) {
  PADDLE_ENFORCE_GE(
      value, 1, platform::errors::InvalidArgument(
                    "send_down_per_steps must >= 1, but received %ld", value));
  send_down_per_steps_ = value;
}

108 109
}  // namespace distributed
}  // namespace paddle