FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

GitHub Viewer

/*************************************************************************** * * 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())

Back | FazBrowse Home | New Git URL