dag_view.cpp 5.9 KB
Newer Older
W
wangguibao 已提交
1
#include "framework/dag_view.h"
W
wangguibao 已提交
2
#include <brpc/traceprintf.h> // TRACEPRINTF
W
wangguibao 已提交
3 4 5 6 7 8 9 10 11 12
#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
    _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) {
20
            LOG(ERROR) << "Failed get stage by index:" << si;
W
wangguibao 已提交
21 22
            return ERR_INTERNAL_FAILURE;
        }
W
wangguibao 已提交
23
        ViewStage* vstage = butil::get_object<ViewStage>();
W
wangguibao 已提交
24
        if (vstage == NULL) {
25
            LOG(ERROR) 
W
wangguibao 已提交
26 27 28 29 30 31 32 33 34
                << "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
            if (vnode == NULL) {
37
                LOG(ERROR) << "Failed get vnode at:" << ni;
W
wangguibao 已提交
38 39 40 41 42
                return ERR_MEM_ALLOC_FAILURE;
            }
            // factory type
            Op* op = OpRepository::instance().get_op(node->type);
            if (op == NULL) {
43
                LOG(ERROR) << "Failed get op with type:" 
W
wangguibao 已提交
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
                    << 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
    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) {
94
            LOG(ERROR) 
W
wangguibao 已提交
95 96 97 98 99 100 101 102 103 104 105 106
                << "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
    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) {
118
            LOG(ERROR) 
W
wangguibao 已提交
119 120 121 122 123
                << "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
                << "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) {
159
        LOG(ERROR) << "invalid empty view stage!";
W
wangguibao 已提交
160 161 162 163 164 165
        return NULL;
    }

    ViewStage* last_stage = _view[_view.size() - 1];
    if (last_stage->nodes.size() != 1 
                    || last_stage->nodes[0] == NULL) {
166
        LOG(ERROR) << "Invalid last stage, size[" 
W
wangguibao 已提交
167
                << last_stage->nodes.size()
W
wangguibao 已提交
168
                << "] != 1";
W
wangguibao 已提交
169 170 171 172 173
        return NULL;
    }

    Op* last_op = last_stage->nodes[0]->op;
    if (last_op == NULL) {
174
        LOG(ERROR) << "Last op is NULL";
W
wangguibao 已提交
175 176 177 178 179 180 181 182
        return NULL;
    }
    return last_op->mutable_channel();
}

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