dag_view.cpp 5.9 KB
Newer Older
W
wangguibao 已提交
1 2 3 4 5 6 7 8 9 10 11 12
#include "framework/dag_view.h"
#include <baidu/rpc/traceprintf.h> // TRACEPRINTF
#include "common/inner_common.h"
#include "framework/op_repository.h"

namespace baidu {
namespace paddle_serving {
namespace predictor {

int DagView::init(Dag* dag, const std::string& service_name) {
    _name = dag->name();
    _full_name = service_name + NAME_DELIMITER + dag->name();
W
wangguibao 已提交
13
    _bus = butil::get_object<Bus>();
W
wangguibao 已提交
14 15 16 17 18 19 20 21 22
    _bus->clear();
    uint32_t stage_size = dag->stage_size();
    // create tls stage view
    for (uint32_t si = 0; si < stage_size; si++) {
        const DagStage* stage = dag->stage_by_index(si);
        if (stage == NULL) {
            LOG(FATAL) << "Failed get stage by index:" << si;
            return ERR_INTERNAL_FAILURE;
        }
W
wangguibao 已提交
23
        ViewStage* vstage = butil::get_object<ViewStage>();
W
wangguibao 已提交
24 25 26 27 28 29 30 31 32 33 34
        if (vstage == NULL) {
            LOG(FATAL) 
                << "Failed get vstage from object pool" 
                << "at:" << si;
            return ERR_MEM_ALLOC_FAILURE;
        }
        vstage->full_name = service_name + NAME_DELIMITER + stage->full_name;
        uint32_t node_size = stage->nodes.size();
        // create tls view node
        for (uint32_t ni = 0; ni < node_size; ni++) {
            DagNode* node = stage->nodes[ni];
W
wangguibao 已提交
35
            ViewNode* vnode = butil::get_object<ViewNode>();
W
wangguibao 已提交
36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74
            if (vnode == NULL) {
                LOG(FATAL) << "Failed get vnode at:" << ni;
                return ERR_MEM_ALLOC_FAILURE;
            }
            // factory type
            Op* op = OpRepository::instance().get_op(node->type);
            if (op == NULL) {
                LOG(FATAL) << "Failed get op with type:" 
                    << node->type;
                return ERR_INTERNAL_FAILURE;
            }

            // initialize a TLS op object
            if (op->init(_bus, dag, node->id, node->name, node->type, node->conf) != 0) {
                LOG(WARNING) << "Failed init op, type:" << node->type;
                return ERR_INTERNAL_FAILURE;
            }
            op->set_full_name(service_name + NAME_DELIMITER + node->full_name);
            vnode->conf = node;
            vnode->op = op;
            vstage->nodes.push_back(vnode);
        }
        _view.push_back(vstage);
    }

    return ERR_OK;
}

int DagView::deinit() {
    uint32_t stage_size = _view.size();
    for (uint32_t si = 0; si < stage_size; si++) {
        ViewStage* vstage = _view[si];
        uint32_t node_size = vstage->nodes.size();
        for (uint32_t ni = 0; ni < node_size; ni++) {
            ViewNode* vnode = vstage->nodes[ni];
            vnode->op->deinit();
            OpRepository::instance().return_op(vnode->op);
            vnode->reset();
            // clear item
W
wangguibao 已提交
75
            butil::return_object(vnode);
W
wangguibao 已提交
76 77 78
        }
        // clear vector
        vstage->nodes.clear();
W
wangguibao 已提交
79
        butil::return_object(vstage);
W
wangguibao 已提交
80 81 82
    }
    _view.clear();
    _bus->clear();
W
wangguibao 已提交
83
    butil::return_object(_bus);
W
wangguibao 已提交
84 85 86
    return ERR_OK;
}

W
wangguibao 已提交
87
int DagView::execute(butil::IOBufBuilder* debug_os) {
W
wangguibao 已提交
88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106
    uint32_t stage_size = _view.size();
    for (uint32_t si = 0; si < stage_size; si++) {
        TRACEPRINTF("start to execute stage[%u]", si);
        int errcode = execute_one_stage(_view[si], debug_os);
        TRACEPRINTF("finish to execute stage[%u]", si);
        if (errcode < 0) {
            LOG(FATAL) 
                << "failed execute stage[" 
                << _view[si]->debug();
            return errcode;
        }
    }
    return ERR_OK;
}

// The default execution strategy is in sequencing
// You can derive a subclass to implement this func.
// ParallelDagView maybe the one you want.
int DagView::execute_one_stage(ViewStage* vstage,
W
wangguibao 已提交
107 108
        butil::IOBufBuilder* debug_os) {
    butil::Timer stage_time(butil::Timer::STARTED);
W
wangguibao 已提交
109 110 111 112 113 114 115 116 117 118 119 120 121 122 123
    uint32_t node_size = vstage->nodes.size(); 
    for (uint32_t ni = 0; ni < node_size; ni++) {
        ViewNode* vnode = vstage->nodes[ni];
        DagNode* conf = vnode->conf;
        Op* op = vnode->op;
        TRACEPRINTF("start to execute op[%s]", op->name());
        int errcode = op->process(debug_os != NULL);
        TRACEPRINTF("finish to execute op[%s]", op->name());
        if (errcode < 0) {
            LOG(FATAL) 
                << "Execute failed, Op:" << op->debug_string();
            return errcode;
        }

        if (errcode > 0) {
W
wangguibao 已提交
124
            LOG(INFO) 
W
wangguibao 已提交
125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182
                << "Execute ignore, Op:" << op->debug_string();
            continue;
        }

        if (debug_os) {
            (*debug_os)
                << "{\"op_name\": \"" << op->name()
                << "\", \"debug_str:\": \"" 
                << op->debug_string()
                << "\", \"time_info\": \"" << op->time_info() << "\"}";
        }

        //LOG(DEBUG) << "Execute succ, Op:" << op->debug_string();
    }
    stage_time.stop();
    PredictorMetric::GetInstance()->update_latency_metric(
            STAGE_METRIC_PREFIX + vstage->full_name, stage_time.u_elapsed());
    return ERR_OK;
}

int DagView::set_request_channel(Channel& request) {
    // Each workflow should get the very beginning 
    // request (channel), and commit it to bus, for
    // the first stage ops consuming.

    request.share_to_bus(_bus);

    return ERR_OK;
}

const Channel* DagView::get_response_channel() const {
    // Caller obtains response channel from bus, and
    // writes it to rpc response(protbuf/json)
    if (_view.size() < 1) {
        LOG(FATAL) << "invalid empty view stage!" << noflush;
        return NULL;
    }

    ViewStage* last_stage = _view[_view.size() - 1];
    if (last_stage->nodes.size() != 1 
                    || last_stage->nodes[0] == NULL) {
        LOG(FATAL) << "Invalid last stage, size[" 
                << last_stage->nodes.size()
                << "] != 1" << noflush;
        return NULL;
    }

    Op* last_op = last_stage->nodes[0]->op;
    if (last_op == NULL) {
        LOG(FATAL) << "Last op is NULL";
        return NULL;
    }
    return last_op->mutable_channel();
}

} // predictor
} // paddle_serving
} // baidu