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