/***************************************************************************
*
* 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
#include "boost/python.hpp"
#include "python_client.h"
#include // NOLINT(readability/streams)
#include
#include "flume/flume.h"
#include "flume/planner/local/local_planner.h"
#include "flume/runtime/backend.h"
#include "flume/core/logical_plan.h"
#include "flume/runtime/common/memory_dataset.h"
#include "flume/runtime/local/local_backend.h"
#include "flume/runtime/local/local_executor_factory.h"
#include "flume/runtime/task.h"
#include "flume/util/reflection.h"
#include "boost/shared_ptr.hpp"
#include "boost/python.hpp"
#include "toft/base/scoped_ptr.h"
#include "bigflow_python/register.h"
#include "bigflow_python/proto/python_resource.pb.h"
DEFINE_bool(bigflow_python_keep_resource, false, "do not delete resource");
DEFINE_bool(bigflow_python_test, false, "run python API locally");
DEFINE_bool(bigflow_python_local, false, "run python API locally");
DEFINE_string(bigflow_python_path, ".", "path of logical plan message");
//DEFINE_string(bigflow_python_resource_path, "./resources", "path of resources message");
namespace baidu {
namespace bigflow {
namespace python {
flume::runtime::CounterSession* g_counter_session = new flume::runtime::CounterSession();
namespace {
// a simple impl, just to show the registered items to stdout
class ReflectionRegister {
public:
template
void add(const std::string& key) {
using flume::Reflection;
std::string name = Reflection::template TypeName();
std::cout AddPythonLibrary(
added_egg_file.file_name(),
added_egg_file.file_path());
}
for (int i = 0; i < python_resource.binary_file_size(); ++i) {
const PbAddedBinary& added_binary_file = python_resource.binary_file(i);
resource->AddFileFromBytes(
added_binary_file.file_name(),
added_binary_file.binary().data(),
added_binary_file.binary().size());
}
if (python_resource.has_cache_file_list()) {
resource->SetCacheFileList(python_resource.cache_file_list());
}
if (python_resource.has_cache_archive_list()) {
resource->SetCacheArchiveList(python_resource.cache_archive_list());
}
backend.SetJobCommitArgs(hadoop_commit_args);
flume::core::LogicalPlan::Status result = backend.Launch(
logical_plan,
resource,
g_counter_session);
if (result) {
LOG(INFO)