From 5b19a56347c1785accd83f97136cc026ef57d386 Mon Sep 17 00:00:00 2001 From: Rob Berlang Date: Wed, 30 Oct 2019 16:29:32 +0100 Subject: [PATCH 1/3] Queue batch prediction tasks collectively --- integration-tests/clipper_admin_tests.py | 42 ++++++- src/benchmarks/src/end_to_end_bench.cpp | 28 +++-- src/frontends/src/query_frontend.hpp | 27 ++--- src/frontends/src/query_frontend_tests.cpp | 84 +++++-------- src/libclipper/include/clipper/datatypes.hpp | 25 ++-- .../include/clipper/query_processor.hpp | 5 +- .../include/clipper/selection_policies.hpp | 38 +++--- .../include/clipper/task_executor.hpp | 110 ++++++++++-------- src/libclipper/src/datatypes.cpp | 19 ++- src/libclipper/src/query_processor.cpp | 97 ++++++++------- src/libclipper/src/selection_policies.cpp | 71 +++++++---- src/libclipper/src/task_executor.cpp | 4 +- .../test/selection_policies_test.cpp | 60 +++++----- src/libclipper/test/task_executor_test.cpp | 3 +- 14 files changed, 342 insertions(+), 271 deletions(-) diff --git a/integration-tests/clipper_admin_tests.py b/integration-tests/clipper_admin_tests.py index 148380247..288d5d6c2 100644 --- a/integration-tests/clipper_admin_tests.py +++ b/integration-tests/clipper_admin_tests.py @@ -791,6 +791,45 @@ def predict_func(inputs): self.assertGreaterEqual(num_max_batch_queries, int(total_num_queries * .7)) + def test_fixed_batch_size_model_processes_all_inputs_as_single_batch( + self): + model_version = 1 + + def predict_func(inputs): + batch_size = len(inputs) + return [str(batch_size) for _ in inputs] + + fixed_batch_size = 9 + total_num_queries = fixed_batch_size + deploy_python_closure( + self.clipper_conn, + self.model_name_4, + model_version, + self.input_type, + predict_func, + batch_size=fixed_batch_size) + self.clipper_conn.link_model_to_app(self.app_name_4, self.model_name_4) + time.sleep(60) + + addr = self.clipper_conn.get_query_addr() + url = "http://{addr}/{app}/predict".format( + addr=addr, app=self.app_name_4) + test_input = [[float(x) + (j * .001) for x in range(5)] + for j in range(total_num_queries)] + req_json = json.dumps({'input_batch': test_input}) + headers = {'Content-type': 'application/json'} + response = requests.post(url, headers=headers, data=req_json) + parsed_response = response.json() + num_max_batch_queries = 0 + for prediction in parsed_response["batch_predictions"]: + batch_size = prediction["output"] + if batch_size != self.default_output and int( + batch_size) == fixed_batch_size: + num_max_batch_queries += 1 + + self.assertEqual(num_max_batch_queries, + total_num_queries) + def test_remove_inactive_container(self): container_name = "{}/noop-container:{}".format(clipper_registry, clipper_version) @@ -902,7 +941,8 @@ def test_remove_inactive_container(self): 'test_deployed_model_queried_successfully', 'test_batch_queries_returned_successfully', 'test_deployed_python_closure_queried_successfully', - 'test_fixed_batch_size_model_processes_specified_query_batch_size_when_saturated' + 'test_fixed_batch_size_model_processes_specified_query_batch_size_when_saturated', + 'test_fixed_batch_size_model_processes_all_inputs_as_single_batch' ] if __name__ == '__main__': diff --git a/src/benchmarks/src/end_to_end_bench.cpp b/src/benchmarks/src/end_to_end_bench.cpp index 4dfeaa93a..a3d0d6cd0 100644 --- a/src/benchmarks/src/end_to_end_bench.cpp +++ b/src/benchmarks/src/end_to_end_bench.cpp @@ -80,28 +80,32 @@ void send_predictions(std::unordered_map &config, query_data_raw[1] = thread_id; } - std::shared_ptr input = std::make_shared( - std::move(query_data), query_vec.size()); + std::vector> input_batch; + input_batch.push_back(std::make_shared( + std::move(query_data), query_vec.size())); Query q = {TEST_APPLICATION_LABEL, UID, - input, + input_batch, latency_objective, clipper::DefaultOutputSelectionPolicy::get_name(), {VersionedModelId(model_name, model_version)}}; - folly::Future prediction = qp.predict(q); + folly::Future> prediction = qp.predict(q); bench_metrics.request_throughput_->mark(1); - std::move(prediction).thenValue([bench_metrics](Response r) { + std::move(prediction) + .thenValue([bench_metrics](std::vector responses) { // Update metrics - if (r.output_is_default_) { - bench_metrics.default_pred_ratio_->increment(1, 1); - } else { - bench_metrics.default_pred_ratio_->increment(0, 1); + for (const auto &r : responses) { + if (r.output_is_default_) { + bench_metrics.default_pred_ratio_->increment(1, 1); + } else { + bench_metrics.default_pred_ratio_->increment(0, 1); + } + bench_metrics.latency_->insert(r.duration_micros_); + bench_metrics.num_predictions_->increment(1); + bench_metrics.throughput_->mark(1); } - bench_metrics.latency_->insert(r.duration_micros_); - bench_metrics.num_predictions_->increment(1); - bench_metrics.throughput_->mark(1); }); } diff --git a/src/frontends/src/query_frontend.hpp b/src/frontends/src/query_frontend.hpp index e25e3f98d..b4d88ee60 100644 --- a/src/frontends/src/query_frontend.hpp +++ b/src/frontends/src/query_frontend.hpp @@ -346,17 +346,16 @@ class RequestHandler { std::shared_ptr response, std::shared_ptr request) { try { - folly::Future>> predictions = + folly::Future> predictions = decode_and_handle_predict(request->content.string(), name, policy, latency_slo_micros, input_type); std::move(predictions) .thenValue([response, - app_metrics](std::vector> tries) { + app_metrics](std::vector responses) { std::vector all_content; - for (auto t : tries) { + for (const auto& r : responses) { try { - Response r = t.value(); if (r.output_is_default_) { app_metrics.default_pred_ratio_->increment(1, 1); } else { @@ -458,10 +457,10 @@ class RequestHandler { } static const std::string parse_output_y_hat( - std::shared_ptr& y_hat) { + const std::shared_ptr& y_hat) { SharedPoolPtr str_content = clipper::get_data(y_hat); - return std::string(str_content.get() + y_hat->start(), - str_content.get() + y_hat->start() + y_hat->size()); + return std::string(str_content.get() + y_hat->start_byte(), + str_content.get() + y_hat->start_byte() + y_hat->byte_size()); } void delete_application(std::string name) { @@ -483,7 +482,7 @@ class RequestHandler { * } */ static const std::string get_prediction_response_content( - Response& query_response) { + const Response& query_response) { rapidjson::Document json_response; json_response.SetObject(); clipper::json::add_long(json_response, PREDICTION_RESPONSE_KEY_QUERY_ID, @@ -570,7 +569,7 @@ class RequestHandler { return clipper::json::to_json_string(error_response); } - folly::Future>> decode_and_handle_predict( + folly::Future> decode_and_handle_predict( std::string content, std::string name, std::string policy, long latency_slo_micros, InputType input_type) { rapidjson::Document d; @@ -626,13 +625,9 @@ class RequestHandler { std::vector> input_batch = clipper::json::parse_inputs(input_type, d); - std::vector> predictions; - for (auto input : input_batch) { - auto prediction = query_processor_.predict(Query{ - name, uid, input, latency_slo_micros, policy, versioned_models}); - predictions.push_back(std::move(prediction)); - } - return folly::collectAll(predictions); + auto predictions = query_processor_.predict(Query{ + name, uid, input_batch, latency_slo_micros, policy, versioned_models}); + return predictions; } /* diff --git a/src/frontends/src/query_frontend_tests.cpp b/src/frontends/src/query_frontend_tests.cpp index 5b27aa1f8..1dba545b2 100644 --- a/src/frontends/src/query_frontend_tests.cpp +++ b/src/frontends/src/query_frontend_tests.cpp @@ -20,12 +20,17 @@ namespace { class MockQueryProcessor { public: MockQueryProcessor() = default; - folly::Future predict(Query query) { - Response response(query, 3, 5, Output("-1.0", {VersionedModelId("m", "1")}), - false, boost::optional{}); - return folly::makeFuture(response); + folly::Future> predict(const Query& query) { + std::vector responses; + responses.reserve(query.input_batch_.size()); + for (const auto& input : query.input_batch_) { + responses.emplace_back( + 3, 5, Output(input, {VersionedModelId("m", "1")}), false, + boost::optional{}); + } + return folly::makeFuture(responses); } - folly::Future update(FeedbackQuery /*feedback*/) { + folly::Future update(const FeedbackQuery& /*feedback*/) { return folly::makeFuture(true); } @@ -68,25 +73,20 @@ class QueryFrontendTest : public ::testing::Test { TEST_F(QueryFrontendTest, TestDecodeCorrectInputInts) { std::string test_json_ints = "{\"input\": [1,2,3,4]}"; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_ints, "test", "test_policy", 30000, InputType::Ints) .get(); - Response response = responses[0].value(); + const Response& response = responses[0]; - Query parsed_query = response.query_; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); int* data = get_data(parsed_input).get(); std::vector parsed_input_data( data + parsed_input->start(), data + parsed_input->start() + parsed_input->size()); - std::vector expected_input_data{1, 2, 3, 4}; EXPECT_EQ(parsed_input_data, expected_input_data); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } TEST_F(QueryFrontendTest, TestDecodeCorrectInputIntsBatch) { @@ -94,40 +94,34 @@ TEST_F(QueryFrontendTest, TestDecodeCorrectInputIntsBatch) { "{\"input_batch\": [[1, 2], [10, 20], [100, 200]]}"; std::vector> expected_input_data{ {1, 2}, {10, 20}, {100, 200}}; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_ints, "test", "test_policy", 30000, InputType::Ints) .get(); for (size_t index = 0; index < responses.size(); ++index) { - Response response = responses[index].value(); - Query parsed_query = response.query_; + const Response& response = responses[index]; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); int* data = get_data(parsed_input).get(); std::vector parsed_input_data( data + parsed_input->start(), data + parsed_input->start() + parsed_input->size()); EXPECT_EQ(parsed_input_data, expected_input_data[index]); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } } TEST_F(QueryFrontendTest, TestDecodeCorrectInputDoubles) { std::string test_json_doubles = "{\"input\": [1.4,2.23,3.243242,0.3223424]}"; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_doubles, "test", "test_policy", 30000, InputType::Doubles) .get(); - Response response = responses[0].value(); - - Query parsed_query = response.query_; + const Response& response = responses[0]; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); double* data = get_data(parsed_input).get(); std::vector parsed_input_data( data + parsed_input->start(), @@ -135,9 +129,6 @@ TEST_F(QueryFrontendTest, TestDecodeCorrectInputDoubles) { std::vector expected_input_data{1.4, 2.23, 3.243242, 0.3223424}; EXPECT_EQ(parsed_input_data, expected_input_data); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } TEST_F(QueryFrontendTest, TestDecodeCorrectInputDoublesBatch) { @@ -145,25 +136,21 @@ TEST_F(QueryFrontendTest, TestDecodeCorrectInputDoublesBatch) { "{\"input_batch\": [[1.1, 2.2], [10.1, 20.2], [100.1, 200.2]]}"; std::vector> expected_input_data{ {1.1, 2.2}, {10.1, 20.2}, {100.1, 200.2}}; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_doubles, "test", "test_policy", 30000, InputType::Doubles) .get(); for (size_t index = 0; index < responses.size(); ++index) { - Response response = responses[index].value(); - Query parsed_query = response.query_; + const Response& response = responses[index]; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); double* data = get_data(parsed_input).get(); std::vector parsed_input_data( data + parsed_input->start(), data + parsed_input->start() + parsed_input->size()); EXPECT_EQ(parsed_input_data, expected_input_data[index]); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } } @@ -171,16 +158,14 @@ TEST_F(QueryFrontendTest, TestDecodeCorrectInputString) { std::string test_json_string = "{\"input\": \"hello world. This is a test string with " "punctionation!@#$Y#;}#\"}"; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_string, "test", "test_policy", 30000, InputType::Strings) .get(); - Response response = responses[0].value(); - - Query parsed_query = response.query_; + const Response& response = responses[0]; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); char* data = get_data(parsed_input).get(); std::string parsed_input_data( data + parsed_input->start(), @@ -189,34 +174,27 @@ TEST_F(QueryFrontendTest, TestDecodeCorrectInputString) { std::string expected_input_data( "hello world. This is a test string with punctionation!@#$Y#;}#"); EXPECT_EQ(parsed_input_data, expected_input_data); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } TEST_F(QueryFrontendTest, TestDecodeCorrectInputStringBatch) { std::string test_json_strings = "{\"input_batch\": [ \"this\", \"is\", \"a\", \"test\" ]}"; std::vector expected_input_data{"this", "is", "a", "test"}; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json_strings, "test", "test_policy", 30000, InputType::Strings) .get(); for (size_t index = 0; index < responses.size(); ++index) { - Response response = responses[index].value(); - Query parsed_query = response.query_; + const Response& response = responses[index]; std::shared_ptr parsed_input = - std::dynamic_pointer_cast(parsed_query.input_); + std::dynamic_pointer_cast(response.output_.y_hat_); char* data = get_data(parsed_input).get(); std::string parsed_input_data( data + parsed_input->start(), data + parsed_input->start() + parsed_input->size()); EXPECT_EQ(parsed_input_data, expected_input_data[index]); - EXPECT_EQ(parsed_query.label_, "test"); - EXPECT_EQ(parsed_query.latency_budget_micros_, 30000); - EXPECT_EQ(parsed_query.selection_policy_, "test_policy"); } } @@ -339,11 +317,11 @@ TEST_F(QueryFrontendTest, TestDeleteManyApplications) { TEST_F(QueryFrontendTest, TestJsonResponseForSuccessfulPredictionFormattedCorrectly) { std::string test_json = "{\"uid\": 1, \"input\": [1,2,3]}"; - std::vector> responses = + std::vector responses = rh_.decode_and_handle_predict(test_json, "test", "test_policy", 30000, InputType::Ints) .get(); - Response response = responses[0].value(); + const Response& response = responses[0]; std::string json_response = rh_.get_prediction_response_content(response); rapidjson::Document parsed_response; @@ -360,7 +338,7 @@ TEST_F(QueryFrontendTest, ->value.IsInt()); ASSERT_TRUE(parsed_response.GetObject() .FindMember(PREDICTION_RESPONSE_KEY_OUTPUT) - ->value.IsFloat()); + ->value.IsString()); ASSERT_TRUE(parsed_response.GetObject() .FindMember(PREDICTION_RESPONSE_KEY_USED_DEFAULT) ->value.IsBool()); diff --git a/src/libclipper/include/clipper/datatypes.hpp b/src/libclipper/include/clipper/datatypes.hpp index 3fadf3e5d..1a44071ab 100644 --- a/src/libclipper/include/clipper/datatypes.hpp +++ b/src/libclipper/include/clipper/datatypes.hpp @@ -76,7 +76,7 @@ class PredictionData { virtual DataType type() const = 0; - virtual PredictionDataHash hash() = 0; + virtual PredictionDataHash hash() const = 0; /** * The index marking the beginning of the input data @@ -177,9 +177,9 @@ class DataVector : public PredictionData { DataType type() const override { return VectorDataType::type; } - PredictionDataHash hash() override { + PredictionDataHash hash() const override { if (!hash_) { - hash_ = CityHash64(reinterpret_cast(data_.get() + start_), + hash_ = CityHash64(reinterpret_cast(data_.get() + start_), size_ * sizeof(D)); } return hash_.get(); @@ -209,7 +209,7 @@ class DataVector : public PredictionData { SharedPoolPtr data_; size_t start_; size_t size_; - boost::optional hash_; + mutable boost::optional hash_; }; typedef DataVector ByteVector; @@ -225,7 +225,8 @@ class Query { public: ~Query() = default; - Query(std::string label, long user_id, std::shared_ptr input, + Query(std::string label, long user_id, + std::vector> input_batch, long latency_budget_micros, std::string selection_policy, std::vector candidate_models); @@ -244,7 +245,7 @@ class Query { // REST endpoints. std::string label_; long user_id_; - std::shared_ptr input_; + std::vector> input_batch_; // TODO change this to a deadline instead of a duration long latency_budget_micros_; std::string selection_policy_; @@ -280,7 +281,7 @@ class Response { public: ~Response() = default; - Response(Query query, QueryId query_id, const long duration_micros, + Response(QueryId query_id, const long duration_micros, Output output, const bool is_default, const boost::optional default_explanation); @@ -294,7 +295,6 @@ class Response { std::string debug_string() const noexcept; - Query query_; QueryId query_id_; long duration_micros_; Output output_; @@ -344,8 +344,8 @@ class PredictTask { public: ~PredictTask() = default; - PredictTask(std::shared_ptr input, VersionedModelId model, - float utility, QueryId query_id, long latency_slo_micros, + PredictTask(std::shared_ptr input, + QueryId query_id, long latency_slo_micros, bool artificial = false); PredictTask(const PredictTask &other) = default; @@ -357,8 +357,6 @@ class PredictTask { PredictTask &operator=(PredictTask &&other) = default; std::shared_ptr input_; - VersionedModelId model_; - float utility_; QueryId query_id_; long latency_slo_micros_; std::chrono::time_point recv_time_; @@ -371,7 +369,7 @@ class FeedbackTask { public: ~FeedbackTask() = default; - FeedbackTask(Feedback feedback, VersionedModelId model, QueryId query_id, + FeedbackTask(Feedback feedback, QueryId query_id, long latency_slo_micros); FeedbackTask(const FeedbackTask &other) = default; @@ -383,7 +381,6 @@ class FeedbackTask { FeedbackTask &operator=(FeedbackTask &&other) = default; Feedback feedback_; - VersionedModelId model_; QueryId query_id_; long latency_slo_micros_; }; diff --git a/src/libclipper/include/clipper/query_processor.hpp b/src/libclipper/include/clipper/query_processor.hpp index 03b8e0503..659c6ad71 100644 --- a/src/libclipper/include/clipper/query_processor.hpp +++ b/src/libclipper/include/clipper/query_processor.hpp @@ -5,6 +5,7 @@ #include #include #include +#include #include @@ -34,8 +35,8 @@ class QueryProcessor { QueryProcessor(QueryProcessor&& other) = default; QueryProcessor& operator=(QueryProcessor&& other) = default; - folly::Future predict(Query query); - folly::Future update(FeedbackQuery feedback); + folly::Future> predict(const Query& query); + folly::Future update(const FeedbackQuery& feedback); std::shared_ptr get_state_table() const; diff --git a/src/libclipper/include/clipper/selection_policies.hpp b/src/libclipper/include/clipper/selection_policies.hpp index cc2b8d503..cf224a575 100644 --- a/src/libclipper/include/clipper/selection_policies.hpp +++ b/src/libclipper/include/clipper/selection_policies.hpp @@ -4,6 +4,7 @@ #include #include #include +#include #include "datatypes.hpp" #include "task_executor.hpp" @@ -53,17 +54,17 @@ class SelectionPolicy { virtual ~SelectionPolicy() = default; // Query Pre-processing: select models and generate tasks - virtual std::vector select_predict_tasks( - std::shared_ptr state, Query query, - long query_id) const = 0; + virtual std::pair, std::vector> + select_predict_tasks(const std::shared_ptr& state, + const Query& query, long first_subquery_id) const = 0; /// Combines multiple prediction results to produce a single - /// output for a query. + /// output for each subquery. /// - /// @returns A pair containing the output and a boolean flag - /// indicating whether or not it is the default output - virtual const std::pair combine_predictions( - const std::shared_ptr& state, Query query, + /// @returns A vector of pairs containing the output and a boolean flag + /// indicating whether or not it is the default output for each subquery. + virtual std::vector> combine_predictions( + const std::shared_ptr& state, const Query& query, std::vector predictions) const = 0; /// When feedback is received, the selection policy can choose @@ -71,9 +72,10 @@ class SelectionPolicy { /// can be used to get y_hat for e.g. updating a bandit algorithm, /// while feedback tasks can be used to optionally propogate feedback /// into the model containers. - virtual std::pair, std::vector> + virtual std::tuple, std::vector, + std::vector> select_feedback_tasks(const std::shared_ptr& state, - FeedbackQuery query, long query_id) const = 0; + const FeedbackQuery& query, long query_id) const = 0; /// This method will be called if at least one PredictTask /// was scheduled for this piece of feedback. This method @@ -128,17 +130,19 @@ class DefaultOutputSelectionPolicy : public SelectionPolicy { std::shared_ptr init_state(Output default_output) const; - std::vector select_predict_tasks( - std::shared_ptr state, Query query, - long query_id) const override; + std::pair, std::vector> + select_predict_tasks(const std::shared_ptr& state, + const Query& query, + long first_subquery_id) const override; - const std::pair combine_predictions( - const std::shared_ptr& state, Query query, + std::vector> combine_predictions( + const std::shared_ptr& state, const Query& query, std::vector predictions) const override; - std::pair, std::vector> + std::tuple, std::vector, + std::vector> select_feedback_tasks(const std::shared_ptr& state, - FeedbackQuery query, long query_id) const override; + const FeedbackQuery& query, long query_id) const override; std::shared_ptr process_feedback( std::shared_ptr state, Feedback feedback, diff --git a/src/libclipper/include/clipper/task_executor.hpp b/src/libclipper/include/clipper/task_executor.hpp index 6754485b8..33acb6a44 100644 --- a/src/libclipper/include/clipper/task_executor.hpp +++ b/src/libclipper/include/clipper/task_executor.hpp @@ -90,10 +90,10 @@ class PredictionCache { public: PredictionCache(size_t size_bytes); folly::Future fetch(const VersionedModelId &model, - std::shared_ptr &input); + const std::shared_ptr &input); void put(const VersionedModelId &model, - std::shared_ptr &input, const Output &output); + const std::shared_ptr &input, const Output &output); private: size_t hash(const VersionedModelId &model, size_t input_hash) const; @@ -132,15 +132,25 @@ class ModelQueue { ~ModelQueue() = default; - void add_task(PredictTask task) { - if (!valid_) return; + void add_tasks(std::vector tasks) { + if (!valid_ || tasks.empty()) return; std::lock_guard lock(queue_mutex_); - Deadline deadline = std::chrono::system_clock::now() + - std::chrono::microseconds(task.latency_slo_micros_); - queue_.emplace(deadline, std::move(task)); + for (auto &task : tasks) { + Deadline deadline = std::chrono::system_clock::now() + + std::chrono::microseconds(task.latency_slo_micros_); + queue_.emplace(deadline, std::move(task)); + } + tasks.clear(); queue_not_empty_condition_.notify_one(); } + void add_task(PredictTask task) { + if (!valid_) return; + std::vector tasks; + tasks.push_back(std::move(task)); + add_tasks(tasks); + } + int get_size() { std::unique_lock l(queue_mutex_); return queue_.size(); @@ -470,60 +480,65 @@ class TaskExecutor { TaskExecutor &operator=(TaskExecutor &&other) = default; std::vector> schedule_predictions( - std::vector tasks) { - predictions_counter_->increment(tasks.size()); + const std::vector &tasks, + const std::vector &models) { + predictions_counter_->increment(models.size() * tasks.size()); std::vector> output_futures; - for (auto t : tasks) { + for (const auto &m : models) { // add each task to the queue corresponding to its associated model boost::shared_lock lock(model_queues_mutex_); - auto model_queue_entry = model_queues_.find(t.model_); + auto model_queue_entry = model_queues_.find(m); if (model_queue_entry != model_queues_.end()) { - auto cache_result = cache_->fetch(t.model_, t.input_); - - if (cache_result.isReady()) { - output_futures.push_back(std::move(cache_result)); - boost::shared_lock model_metrics_lock( - model_metrics_mutex_); - auto cur_model_metric_entry = model_metrics_.find(t.model_); - if (cur_model_metric_entry != model_metrics_.end()) { - auto cur_model_metric = cur_model_metric_entry->second; - cur_model_metric.cache_hit_ratio_->increment(1, 1); - } - } - - else if (active_containers_->get_replicas_for_model(t.model_).size() == - 0) { - log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, - "No active model containers for model: {} : {}", - t.model_.get_name(), t.model_.get_id()); - } else { - output_futures.push_back(std::move(cache_result)); - t.recv_time_ = std::chrono::system_clock::now(); - model_queue_entry->second->add_task(t); - log_info_formatted(LOGGING_TAG_TASK_EXECUTOR, - "Adding task to queue. QueryID: {}, model: {}", - t.query_id_, t.model_.serialize()); - boost::shared_lock model_metrics_lock( - model_metrics_mutex_); - auto cur_model_metric_entry = model_metrics_.find(t.model_); - if (cur_model_metric_entry != model_metrics_.end()) { - auto cur_model_metric = cur_model_metric_entry->second; - cur_model_metric.cache_hit_ratio_->increment(0, 1); + std::vector model_tasks; + for (const auto &t : tasks) { + auto cache_result = cache_->fetch(m, t.input_); + + if (cache_result.isReady()) { + output_futures.push_back(std::move(cache_result)); + boost::shared_lock model_metrics_lock( + model_metrics_mutex_); + auto cur_model_metric_entry = model_metrics_.find(m); + if (cur_model_metric_entry != model_metrics_.end()) { + auto cur_model_metric = cur_model_metric_entry->second; + cur_model_metric.cache_hit_ratio_->increment(1, 1); + } + } else if (active_containers_->get_replicas_for_model(m).size() == + 0) { + log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, + "No active model containers for model: {} : {}", + m.get_name(), m.get_id()); + } else { + output_futures.push_back(std::move(cache_result)); + model_tasks.push_back(t); + model_tasks.back().recv_time_ = std::chrono::system_clock::now(); + log_info_formatted(LOGGING_TAG_TASK_EXECUTOR, + "Adding task to queue. QueryID: {}, model: {}", + t.query_id_, m.serialize()); + boost::shared_lock model_metrics_lock( + model_metrics_mutex_); + auto cur_model_metric_entry = model_metrics_.find(m); + if (cur_model_metric_entry != model_metrics_.end()) { + auto cur_model_metric = cur_model_metric_entry->second; + cur_model_metric.cache_hit_ratio_->increment(0, 1); + } } } + model_queue_entry->second->add_tasks(model_tasks); } else { log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, "Received task for unknown model: {} : {}", - t.model_.get_name(), t.model_.get_id()); + m.get_name(), m.get_id()); } } return output_futures; } std::vector> schedule_feedback( - const std::vector tasks) { + const std::vector &tasks, + const std::vector &models) { UNUSED(tasks); + UNUSED(models); // TODO Implement return {}; } @@ -717,10 +732,11 @@ class TaskExecutor { std::stringstream query_ids_in_batch; std::chrono::time_point current_time = std::chrono::system_clock::now(); - for (auto b : batch) { + for (const auto &b : batch) { prediction_request.add_input(b.input_); - cur_batch.emplace_back(current_time, container->container_id_, b.model_, - container->replica_id_, b.input_, b.artificial_); + cur_batch.emplace_back(current_time, container->container_id_, + container->model_, container->replica_id_, + b.input_, b.artificial_); query_ids_in_batch << b.query_id_ << " "; } int message_id = rpc_->send_message(prediction_request.serialize(), diff --git a/src/libclipper/src/datatypes.cpp b/src/libclipper/src/datatypes.cpp index a87a4d094..8383e13c5 100644 --- a/src/libclipper/src/datatypes.cpp +++ b/src/libclipper/src/datatypes.cpp @@ -223,22 +223,21 @@ rpc::PredictionResponse::deserialize_prediction_response( } Query::Query(std::string label, long user_id, - std::shared_ptr input, long latency_budget_micros, - std::string selection_policy, + std::vector> input_batch, + long latency_budget_micros, std::string selection_policy, std::vector candidate_models) : label_(std::move(label)), user_id_(user_id), - input_(std::move(input)), + input_batch_(std::move(input_batch)), latency_budget_micros_(latency_budget_micros), selection_policy_(std::move(selection_policy)), candidate_models_(std::move(candidate_models)), create_time_(std::chrono::high_resolution_clock::now()) {} -Response::Response(Query query, QueryId query_id, const long duration_micros, +Response::Response(QueryId query_id, const long duration_micros, Output output, const bool output_is_default, const boost::optional default_explanation) - : query_(std::move(query)), - query_id_(query_id), + : query_id_(query_id), duration_micros_(duration_micros), output_(std::move(output)), output_is_default_(output_is_default), @@ -266,20 +265,16 @@ FeedbackQuery::FeedbackQuery(std::string label, long user_id, Feedback feedback, candidate_models_(std::move(candidate_models)) {} PredictTask::PredictTask(std::shared_ptr input, - VersionedModelId model, float utility, QueryId query_id, long latency_slo_micros, bool artificial) : input_(std::move(input)), - model_(std::move(model)), - utility_(utility), query_id_(query_id), latency_slo_micros_(latency_slo_micros), artificial_(artificial) {} -FeedbackTask::FeedbackTask(Feedback feedback, VersionedModelId model, - QueryId query_id, long latency_slo_micros) +FeedbackTask::FeedbackTask(Feedback feedback, QueryId query_id, + long latency_slo_micros) : feedback_(feedback), - model_(model), query_id_(query_id), latency_slo_micros_(latency_slo_micros) {} diff --git a/src/libclipper/src/query_processor.cpp b/src/libclipper/src/query_processor.cpp index b43420777..9e0e4d2ca 100644 --- a/src/libclipper/src/query_processor.cpp +++ b/src/libclipper/src/query_processor.cpp @@ -42,8 +42,9 @@ std::shared_ptr QueryProcessor::get_state_table() const { return state_db_; } -folly::Future QueryProcessor::predict(Query query) { - long query_id = query_counter_.fetch_add(1); +folly::Future> QueryProcessor::predict( + const Query &query) { + long first_subquery_id = query_counter_.fetch_add(query.input_batch_.size()); auto current_policy_iter = selection_policies_.find(query.selection_policy_); if (current_policy_iter == selection_policies_.end()) { std::stringstream err_msg_builder; @@ -68,38 +69,39 @@ folly::Future QueryProcessor::predict(Query query) { current_policy->deserialize(*state_opt); boost::optional default_explanation; - std::vector tasks = - current_policy->select_predict_tasks(selection_state, query, query_id); + std::vector tasks; + std::vector models; + std::tie(tasks, models) = current_policy->select_predict_tasks( + selection_state, query, first_subquery_id); log_info_formatted(LOGGING_TAG_QUERY_PROCESSOR, "Found {} tasks", - tasks.size()); + models.size() * tasks.size()); - vector> task_futures = - task_executor_.schedule_predictions(tasks); + std::vector> task_futures = + task_executor_.schedule_predictions(tasks, models); if (task_futures.empty()) { default_explanation = "No connected models found for query"; - log_error_formatted(LOGGING_TAG_QUERY_PROCESSOR, - "No connected models found for query with id: {}", - query_id); + log_error_formatted( + LOGGING_TAG_QUERY_PROCESSOR, + "No connected models found for query with ids: {} to {}", + first_subquery_id, first_subquery_id + query.input_batch_.size() - 1); } - size_t num_tasks = task_futures.size(); - folly::Future timer_future = timer_system_.set_timer(query.latency_budget_micros_); std::shared_ptr outputs_mutex = std::make_shared(); - std::vector outputs; - outputs.reserve(task_futures.size()); std::shared_ptr> outputs_ptr = - std::make_shared>(std::move(outputs)); + std::make_shared>(task_futures.size()); std::vector> wrapped_task_futures; - for (auto it = task_futures.begin(); it < task_futures.end(); it++) { + for (size_t i = 0; i < task_futures.size(); ++i) { + auto &t = task_futures[i]; wrapped_task_futures.push_back( - std::move(*it).thenValue([outputs_mutex, outputs_ptr](Output output) { + std::move(t) + .thenValue([outputs_mutex, outputs_ptr, i](Output output) { std::lock_guard lock(*outputs_mutex); - outputs_ptr->push_back(output); + (*outputs_ptr)[i] = output; }).thenError(folly::tag_t{}, [](const std::exception& e) { log_error_formatted( LOGGING_TAG_QUERY_PROCESSOR, @@ -119,23 +121,20 @@ folly::Future QueryProcessor::predict(Query query) { folly::Future>> response_ready_future = folly::collectAny(when_either_futures); - folly::Promise response_promise; - folly::Future response_future = response_promise.getFuture(); + folly::Promise> response_promise; + folly::Future> response_future = response_promise.getFuture(); - std::move(response_ready_future).thenValue([ - outputs_ptr, outputs_mutex, num_tasks, query, query_id, selection_state, - current_policy, response_promise = std::move(response_promise), - default_explanation + std::move(response_ready_future) + .thenValue( + [outputs_ptr, outputs_mutex, query, first_subquery_id, + selection_state, current_policy, + response_promise = std::move(response_promise), + default_explanation ](const std::pair>& /* completed_future */) mutable { std::lock_guard outputs_lock(*outputs_mutex); - if (outputs_ptr->empty() && num_tasks > 0 && !default_explanation) { - default_explanation = - "Failed to retrieve a prediction response within the specified " - "latency SLO"; - } - std::pair final_output = current_policy->combine_predictions( + std::vector> final_output = current_policy->combine_predictions( selection_state, query, *outputs_ptr); std::chrono::time_point end = @@ -145,18 +144,26 @@ folly::Future QueryProcessor::predict(Query query) { end - query.create_time_) .count(); - Response response{query, - query_id, - duration_micros, - final_output.first, - final_output.second, - default_explanation}; - response_promise.setValue(response); + std::vector responses; + responses.reserve(final_output.size()); + long subquery_id = first_subquery_id; + for (const auto& f : final_output) { + responses.emplace_back( + subquery_id, duration_micros, f.first, f.second, + ((default_explanation || !f.second) + ? default_explanation + : boost::optional("Failed to retrieve a prediction " + "response within the specified " + "latency SLO"))); + ++subquery_id; + } + + response_promise.setValue(responses); }); return response_future; } -folly::Future QueryProcessor::update(FeedbackQuery feedback) { +folly::Future QueryProcessor::update(const FeedbackQuery& feedback) { log_info(LOGGING_TAG_QUERY_PROCESSOR, "Received feedback for user {}", feedback.user_id_); @@ -188,24 +195,26 @@ folly::Future QueryProcessor::update(FeedbackQuery feedback) { std::vector predict_tasks; std::vector feedback_tasks; - std::tie(predict_tasks, feedback_tasks) = + std::vector models; + std::tie(predict_tasks, feedback_tasks, models) = current_policy->select_feedback_tasks(selection_state, feedback, query_id); log_info_formatted(LOGGING_TAG_QUERY_PROCESSOR, "Scheduling {} prediction tasks and {} feedback tasks", - predict_tasks.size(), feedback_tasks.size()); + predict_tasks.size() * models.size(), + feedback_tasks.size() * models.size()); // 1) Wait for all prediction_tasks to complete // 2) Update selection policy // 3) Complete select_policy_update_promise // 4) Wait for all feedback_tasks to complete (feedback_processed future) - vector> predict_task_futures = - task_executor_.schedule_predictions({predict_tasks}); + std::vector> predict_task_futures = + task_executor_.schedule_predictions(predict_tasks, models); - vector> feedback_task_futures = - task_executor_.schedule_feedback(std::move(feedback_tasks)); + std::vector> feedback_task_futures = + task_executor_.schedule_feedback(feedback_tasks, models); folly::Future> all_preds_completed = folly::collect(predict_task_futures); diff --git a/src/libclipper/src/selection_policies.cpp b/src/libclipper/src/selection_policies.cpp index e28ded2e9..61e399cbc 100644 --- a/src/libclipper/src/selection_policies.cpp +++ b/src/libclipper/src/selection_policies.cpp @@ -62,10 +62,11 @@ std::shared_ptr DefaultOutputSelectionPolicy::init_state( return std::make_shared(default_output); } -std::vector DefaultOutputSelectionPolicy::select_predict_tasks( - std::shared_ptr /*state*/, Query query, - long query_id) const { - std::vector tasks; +std::pair, std::vector> +DefaultOutputSelectionPolicy::select_predict_tasks( + const std::shared_ptr& /*state*/, const Query& query, + long first_subquery_id) const { + std::vector models; size_t num_candidate_models = query.candidate_models_.size(); if (num_candidate_models == (size_t)0) { log_error_formatted(LOGGING_TAG_SELECTION_POLICY, @@ -78,37 +79,61 @@ std::vector DefaultOutputSelectionPolicy::select_predict_tasks( "{}. Picking the first one.", num_candidate_models, query.label_); } - tasks.emplace_back(query.input_, query.candidate_models_.front(), 1.0, - query_id, query.latency_budget_micros_); + models.emplace_back(query.candidate_models_.front()); + } + + std::vector tasks; + tasks.reserve(query.input_batch_.size()); + for (const auto& input : query.input_batch_) { + tasks.emplace_back(input, first_subquery_id++, + query.latency_budget_micros_); } - return tasks; + return std::make_pair(std::move(tasks), std::move(models)); } -const std::pair DefaultOutputSelectionPolicy::combine_predictions( - const std::shared_ptr& state, Query /*query*/, +std::vector> +DefaultOutputSelectionPolicy::combine_predictions( + const std::shared_ptr& state, const Query& query, std::vector predictions) const { - if (predictions.size() == 1) { - return std::make_pair(std::move(predictions.front()), false); - } else if (predictions.empty()) { + std::vector> outputs; + outputs.reserve(query.input_batch_.size()); + for (auto& p : predictions) { + if (outputs.size() >= query.input_batch_.size()) { + break; + } + + if (p.y_hat_) { + outputs.emplace_back(std::move(p), false); + } else { + outputs.emplace_back(std::dynamic_pointer_cast(state) + ->default_output_, true); + } + } + + if (outputs.size() < query.input_batch_.size()) { Output default_output = std::dynamic_pointer_cast(state) ->default_output_; - return std::make_pair(std::move(default_output), true); - } else { + outputs.resize(query.input_batch_.size(), + std::make_pair(default_output, true)); + } else if (predictions.size() > query.input_batch_.size()) { log_error_formatted(LOGGING_TAG_SELECTION_POLICY, - "DefaultOutputSelectionPolicy only expecting 1 " - "output but found {}. Returning the first one.", - predictions.size()); - return std::make_pair(std::move(predictions.front()), false); + "DefaultOutputSelectionPolicy only expecting {} " + "outputs but found {}. Returning the first {}.", + query.input_batch_.size(), predictions.size(), + query.input_batch_.size()); } + + return outputs; } -std::pair, std::vector> +std::tuple, std::vector, + std::vector> DefaultOutputSelectionPolicy::select_feedback_tasks( - const std::shared_ptr& /*state*/, FeedbackQuery /*query*/, - long /*query_id*/) const { - return std::make_pair, std::vector>( - {}, {}); + const std::shared_ptr& /*state*/, + const FeedbackQuery& /*query*/, long /*query_id*/) const { + return std::make_tuple, std::vector, + std::vector>({}, {}, {}); } std::shared_ptr DefaultOutputSelectionPolicy::process_feedback( diff --git a/src/libclipper/src/task_executor.cpp b/src/libclipper/src/task_executor.cpp index 0f36e1444..ef3ab06ff 100644 --- a/src/libclipper/src/task_executor.cpp +++ b/src/libclipper/src/task_executor.cpp @@ -22,7 +22,7 @@ PredictionCache::PredictionCache(size_t size_bytes) } folly::Future PredictionCache::fetch( - const VersionedModelId &model, std::shared_ptr &input) { + const VersionedModelId &model, const std::shared_ptr &input) { std::unique_lock l(m_); auto key = hash(model, input->hash()); auto search = entries_.find(key); @@ -60,7 +60,7 @@ folly::Future PredictionCache::fetch( } void PredictionCache::put(const VersionedModelId &model, - std::shared_ptr &input, + const std::shared_ptr &input, const Output &output) { std::unique_lock l(m_); auto key = hash(model, input->hash()); diff --git a/src/libclipper/test/selection_policies_test.cpp b/src/libclipper/test/selection_policies_test.cpp index 615cf7957..b262c142b 100644 --- a/src/libclipper/test/selection_policies_test.cpp +++ b/src/libclipper/test/selection_policies_test.cpp @@ -25,15 +25,17 @@ class DefaultOutputSelectionPolicyTest : public ::testing::Test { TEST_F(DefaultOutputSelectionPolicyTest, TestSelectPredictTasksZeroCandidateModels) { - Query zero_candidate_models_query{"label", - clipper::DEFAULT_USER_ID, - std::shared_ptr(), - 1000, - DefaultOutputSelectionPolicy::get_name(), - {}}; + Query zero_candidate_models_query{ + "label", + clipper::DEFAULT_USER_ID, + std::vector>( + 1, std::shared_ptr()), + 1000, + DefaultOutputSelectionPolicy::get_name(), + {}}; auto zero_models_tasks = policy_.select_predict_tasks(nullptr, zero_candidate_models_query, 0); - EXPECT_EQ(zero_models_tasks.size(), (size_t)0); + EXPECT_EQ(zero_models_tasks.second.size(), (size_t)0); } TEST_F(DefaultOutputSelectionPolicyTest, @@ -43,14 +45,15 @@ TEST_F(DefaultOutputSelectionPolicyTest, VersionedModelId("simple_svm", "2")}; Query two_candidate_models_query{"label", clipper::DEFAULT_USER_ID, - std::shared_ptr(), + std::vector>( + 1, std::shared_ptr()), 1000, DefaultOutputSelectionPolicy::get_name(), two_models}; auto two_models_tasks = policy_.select_predict_tasks(nullptr, two_candidate_models_query, 0); - EXPECT_EQ(two_models_tasks.size(), (size_t)1); - EXPECT_EQ(two_models_tasks.front().model_, two_models.front()); + ASSERT_EQ(two_models_tasks.second.size(), (size_t)1); + EXPECT_EQ(two_models_tasks.second.front(), two_models.front()); } TEST_F(DefaultOutputSelectionPolicyTest, @@ -59,14 +62,15 @@ TEST_F(DefaultOutputSelectionPolicyTest, VersionedModelId("music_random_features", "1")}; Query one_candidate_model_query{"label", clipper::DEFAULT_USER_ID, - std::shared_ptr(), + std::vector>( + 1, std::shared_ptr()), 1000, DefaultOutputSelectionPolicy::get_name(), one_model}; auto one_model_tasks = policy_.select_predict_tasks(nullptr, one_candidate_model_query, 0); - EXPECT_EQ(one_model_tasks.size(), (size_t)1); - EXPECT_EQ(one_model_tasks.front().model_, one_model.front()); + ASSERT_EQ(one_model_tasks.second.size(), (size_t)1); + EXPECT_EQ(one_model_tasks.second.front(), one_model.front()); } TEST_F(DefaultOutputSelectionPolicyTest, @@ -74,20 +78,23 @@ TEST_F(DefaultOutputSelectionPolicyTest, VersionedModelId m1 = VersionedModelId("music_random_features", "1"); Query one_candidate_model_query{"label", clipper::DEFAULT_USER_ID, - std::shared_ptr(), + std::vector>( + 1, std::shared_ptr()), 1000, DefaultOutputSelectionPolicy::get_name(), {m1}}; auto zero_preds_output = - policy_.combine_predictions(state_, one_candidate_model_query, {}).first; - ASSERT_EQ(zero_preds_output, state_->default_output_); + policy_.combine_predictions(state_, one_candidate_model_query, {}); + ASSERT_EQ(zero_preds_output.size(), (size_t)1); + EXPECT_EQ(zero_preds_output.front().first, state_->default_output_); } TEST_F(DefaultOutputSelectionPolicyTest, TestCombinePredictionsOnePrediction) { VersionedModelId m1 = VersionedModelId("music_random_features", "1"); Query one_candidate_model_query{"label", clipper::DEFAULT_USER_ID, - std::shared_ptr(), + std::vector>( + 1, std::shared_ptr()), 1000, DefaultOutputSelectionPolicy::get_name(), {m1}}; @@ -96,10 +103,10 @@ TEST_F(DefaultOutputSelectionPolicyTest, TestCombinePredictionsOnePrediction) { auto one_pred_output = policy_ .combine_predictions(state_, one_candidate_model_query, - {first_output}) - .first; - ASSERT_EQ(one_pred_output, first_output); - ASSERT_NE(one_pred_output, state_->default_output_); + {first_output}); + ASSERT_EQ(one_pred_output.size(), (size_t)1); + EXPECT_EQ(one_pred_output.front().first, first_output); + EXPECT_NE(one_pred_output.front().first, state_->default_output_); } TEST_F(DefaultOutputSelectionPolicyTest, TestCombinePredictionsTwoPredictions) { @@ -107,7 +114,8 @@ TEST_F(DefaultOutputSelectionPolicyTest, TestCombinePredictionsTwoPredictions) { VersionedModelId m2 = VersionedModelId("simple_svm", "2"); Query two_candidate_models_query{"label", clipper::DEFAULT_USER_ID, - std::shared_ptr(), + std::vector>( + 1, std::shared_ptr()), 1000, DefaultOutputSelectionPolicy::get_name(), {m1, m2}}; @@ -116,10 +124,10 @@ TEST_F(DefaultOutputSelectionPolicyTest, TestCombinePredictionsTwoPredictions) { auto two_preds_output = policy_ .combine_predictions(state_, two_candidate_models_query, - {first_output, second_output}) - .first; - ASSERT_EQ(two_preds_output, first_output); - ASSERT_NE(two_preds_output, state_->default_output_); + {first_output, second_output}); + ASSERT_EQ(two_preds_output.size(), (size_t)1); + EXPECT_EQ(two_preds_output.front().first, first_output); + EXPECT_NE(two_preds_output.front().first, state_->default_output_); } TEST(DefaultOutputSelectionStateTest, Serialization) { diff --git a/src/libclipper/test/task_executor_test.cpp b/src/libclipper/test/task_executor_test.cpp index ca3a3b1c4..28ba26bec 100644 --- a/src/libclipper/test/task_executor_test.cpp +++ b/src/libclipper/test/task_executor_test.cpp @@ -21,10 +21,9 @@ namespace { PredictTask create_predict_task(long query_id, long latency_slo_millis) { UniquePoolPtr data = memory::allocate_unique(1); data.get()[0] = 1.0; - VersionedModelId model_id = VersionedModelId("test", "1"); std::shared_ptr input = std::make_shared(std::move(data), 1); - PredictTask task(input, model_id, 1.0, query_id, latency_slo_millis); + PredictTask task(input, query_id, latency_slo_millis); return task; } From ed48164f94db7cb37ca3ec70346ece0954aaaa0d Mon Sep 17 00:00:00 2001 From: Rob Berlang Date: Sun, 1 Dec 2019 15:14:08 +0100 Subject: [PATCH 2/3] Handle case when no active containers found properly --- src/libclipper/include/clipper/task_executor.hpp | 10 +++++++--- src/libclipper/src/selection_policies.cpp | 4 ++++ 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/src/libclipper/include/clipper/task_executor.hpp b/src/libclipper/include/clipper/task_executor.hpp index 33acb6a44..39135f501 100644 --- a/src/libclipper/include/clipper/task_executor.hpp +++ b/src/libclipper/include/clipper/task_executor.hpp @@ -484,13 +484,14 @@ class TaskExecutor { const std::vector &models) { predictions_counter_->increment(models.size() * tasks.size()); std::vector> output_futures; + boost::shared_lock lock(model_queues_mutex_); for (const auto &m : models) { // add each task to the queue corresponding to its associated model - boost::shared_lock lock(model_queues_mutex_); auto model_queue_entry = model_queues_.find(m); if (model_queue_entry != model_queues_.end()) { std::vector model_tasks; + size_t initial_outputs_size = output_futures.size(); for (const auto &t : tasks) { auto cache_result = cache_->fetch(m, t.input_); @@ -503,11 +504,14 @@ class TaskExecutor { auto cur_model_metric = cur_model_metric_entry->second; cur_model_metric.cache_hit_ratio_->increment(1, 1); } - } else if (active_containers_->get_replicas_for_model(m).size() == - 0) { + } else if (active_containers_->get_replicas_for_model(m).empty()) { log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, "No active model containers for model: {} : {}", m.get_name(), m.get_id()); + model_tasks.clear(); + output_futures.erase(output_futures.begin() + initial_outputs_size, + output_futures.end()); + break; } else { output_futures.push_back(std::move(cache_result)); model_tasks.push_back(t); diff --git a/src/libclipper/src/selection_policies.cpp b/src/libclipper/src/selection_policies.cpp index 61e399cbc..4020de538 100644 --- a/src/libclipper/src/selection_policies.cpp +++ b/src/libclipper/src/selection_policies.cpp @@ -111,6 +111,10 @@ DefaultOutputSelectionPolicy::combine_predictions( } if (outputs.size() < query.input_batch_.size()) { + log_error_formatted(LOGGING_TAG_SELECTION_POLICY, + "DefaultOutputSelectionPolicy expecting {} " + "outputs but found {}. Filling with default.", + query.input_batch_.size(), predictions.size()); Output default_output = std::dynamic_pointer_cast(state) ->default_output_; From 2774441e3297b49067ed07e57d392a9d267e6a56 Mon Sep 17 00:00:00 2001 From: Rob Berlang Date: Sat, 29 Aug 2020 16:43:58 +0200 Subject: [PATCH 3/3] Handle case when no active containers found properly --- src/libclipper/include/clipper/task_executor.hpp | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) diff --git a/src/libclipper/include/clipper/task_executor.hpp b/src/libclipper/include/clipper/task_executor.hpp index 39135f501..bf7d1558a 100644 --- a/src/libclipper/include/clipper/task_executor.hpp +++ b/src/libclipper/include/clipper/task_executor.hpp @@ -486,12 +486,17 @@ class TaskExecutor { std::vector> output_futures; boost::shared_lock lock(model_queues_mutex_); for (const auto &m : models) { + if (active_containers_->get_replicas_for_model(m).empty()) { + log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, + "No active model containers for model: {} : {}", + m.get_name(), m.get_id()); + continue; + } // add each task to the queue corresponding to its associated model auto model_queue_entry = model_queues_.find(m); if (model_queue_entry != model_queues_.end()) { std::vector model_tasks; - size_t initial_outputs_size = output_futures.size(); for (const auto &t : tasks) { auto cache_result = cache_->fetch(m, t.input_); @@ -504,14 +509,6 @@ class TaskExecutor { auto cur_model_metric = cur_model_metric_entry->second; cur_model_metric.cache_hit_ratio_->increment(1, 1); } - } else if (active_containers_->get_replicas_for_model(m).empty()) { - log_error_formatted(LOGGING_TAG_TASK_EXECUTOR, - "No active model containers for model: {} : {}", - m.get_name(), m.get_id()); - model_tasks.clear(); - output_futures.erase(output_futures.begin() + initial_outputs_size, - output_futures.end()); - break; } else { output_futures.push_back(std::move(cache_result)); model_tasks.push_back(t);