[ Web Proxy ]
URL:
Viewing: https://raw.githubusercontent.com/himdd/bigflow/master/bigflow_python/python_client.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 

#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) 

Web Proxy Viewer  |  New URL  |  Original Page