[ Web Proxy ]
URL:
Viewing: https://raw.githubusercontent.com/himdd/bigflow/master/bigflow_python/python_interpreter.cpp [Back]  [Original]

/***************************************************************************
 *
 * Copyright (c) 2015 Baidu, Inc. 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.
 *
 **************************************************************************/
// Author: Wang Cong (wangcong09@baidu.com)
//

#include "bigflow_python/python_interpreter.h"

#include 
#include 

#include "boost/filesystem.hpp"
#include "boost/python/suite/indexing/vector_indexing_suite.hpp"
#include "boost/python/def.hpp"
#include "boost/lexical_cast.hpp"
#include "gflags/gflags.h"
#include "glog/logging.h"

#include "toft/base/string/string_piece.h"
#include "toft/base/string/algorithm.h"

#include "bigflow_python/common/python.h"
#include "bigflow_python/delegators/python_processor_delegator.h"
#include "bigflow_python/processors/processor.h"
#include "bigflow_python/serde/cpickle_serde.h"
#include "flume/core/iterator.h"
#include "flume/runtime/io/io_format.h"
#include "flume/runtime/local/flags.h"

DECLARE_string(flume_backend);
DECLARE_bool(flume_commit);
DECLARE_int32(flume_log_server_index);

namespace baidu {
namespace flume {
namespace runtime {
namespace spark {
DECLARE_string(flume_python_home);
DECLARE_string(flume_application_home);
}  // namespace spark
}  // namespace runtime
}  // namespace flume
}  // namespace baidu

namespace baidu {
namespace bigflow {
namespace python {

namespace {

// Translate FlumeIteratorDelegator::StopIteration to Python StopIteration exception
void translate(const FlumeIteratorDelegator::StopIteration& e) {
    // Use the Python C API to set up an exception object
    PyErr_SetNone(PyExc_StopIteration);
}

void expose() {
    boost::python::class_("StringPiece")
        .def("as_string", &toft::StringPiece::as_string)
        .def("data", &toft::StringPiece::data)
    ;

    boost::python::class_("VecStringPiece")
        .def(boost::python::vector_indexing_suite())
    ;

    boost::python::class_("VecIterator")
        .def(boost::python::vector_indexing_suite())
    ;

    boost::python::register_exception_translator(&translate);
    boost::python::class_("FlumeIterator", boost::python::no_init)
        .def("has_next", &FlumeIteratorDelegator::HasNext)
        .def("next", &FlumeIteratorDelegator::NextValue)
        .def("reset", &FlumeIteratorDelegator::Reset)
        .def("done", &FlumeIteratorDelegator::Done)
    ;

    boost::python::class_("FlumeEmitter")
        .def("emit", &EmitterDelegator::emit)
        .def("done", &EmitterDelegator::done)
    ;

    boost::python::class_("Record")
        .def_readwrite("key", &flume::runtime::Record::key)
        .def_readwrite("value", &flume::runtime::Record::value)
    ;

    boost::python::class_("SideInput", boost::python::no_init)
        .def("__len__", &SideInput::length)
        .def("__iter__", &SideInput::get_iterator)
        .def("__getitem__", &SideInput::get_item)
        .def("as_list", &SideInput::as_list)
    ;
}

}  // namespace

PythonInterpreter::PythonInterpreter(): _handler(NULL) {
    // In the python code, we can check this environment var
    // to judge if it is running on the remote side. (or on the client side).
    CHECK_EQ(0, setenv("__PYTHON_IN_REMOTE_SIDE", "true", 1));
    initialize_interpreter();
    expose();
}

PythonInterpreter::~PythonInterpreter() {
    //if (is_interpreter_initialized()) {
    //    finalize_interpreter();
    //}
}

void PythonInterpreter::set_exception_handler(ExceptionHandlerFunc handler) {
    _handler = handler;
}

const void PythonInterpreter::handle_exception(const char* file,
                                               const char* function,
                                               int line_num) {
    (*_handler)(file, function, line_num);
}

void PythonInterpreter::set_exception_handler_with_error_msg(
    ExceptionHandlerWithMsgFunc handler_with_error_msg) {
    _handler_with_error_msg = handler_with_error_msg;
}

const void PythonInterpreter::handle_exception_with_error_msg(const std::string& err_msg,
                                                              const char* file,
                                                              const char* function,
                                                              int line_num) {
    _handler_with_error_msg(err_msg, file, function, line_num);
}

void PythonInterpreter::initialize_interpreter() {
    try {
        CHECK(!is_interpreter_initialized()) 

Web Proxy Viewer  |  New URL  |  Original Page