diff --git a/tensorflow_serving/batching/BUILD b/tensorflow_serving/batching/BUILD index 0ba47577dc3..860896344d1 100644 --- a/tensorflow_serving/batching/BUILD +++ b/tensorflow_serving/batching/BUILD @@ -170,6 +170,7 @@ cc_library( srcs = ["batching_util.cc"], hdrs = ["batching_util.h"], deps = [ + "@com_google_absl//absl/status", "@com_google_absl//absl/strings", "@com_google_absl//absl/types:span", "@org_tensorflow//tensorflow/core:framework", diff --git a/tensorflow_serving/batching/batching_session.cc b/tensorflow_serving/batching/batching_session.cc index d0a73c68736..38294ffd3a1 100644 --- a/tensorflow_serving/batching/batching_session.cc +++ b/tensorflow_serving/batching/batching_session.cc @@ -513,7 +513,9 @@ absl::Status BatchingSession::MergeInputTensors( if (options_.pad_variable_length_inputs) { std::vector>> all_task_inputs = GetTaskInputsVector(batch); - max_dim_sizes = CalculateMaxDimSizes(all_task_inputs); + max_dim_sizes.emplace(); + TF_RETURN_IF_ERROR( + CalculateMaxDimSizes(all_task_inputs, &max_dim_sizes.value())); } // Populate 'tensors_to_merge'. for (int i = 0; i < batch.num_tasks(); ++i) { diff --git a/tensorflow_serving/batching/batching_session_test.cc b/tensorflow_serving/batching/batching_session_test.cc index 6c4fb9ec4b5..0518ac48d48 100644 --- a/tensorflow_serving/batching/batching_session_test.cc +++ b/tensorflow_serving/batching/batching_session_test.cc @@ -344,6 +344,35 @@ TEST_P(BatchingSessionTest, BatchingWithPadding) { })); } +TEST(BatchingSessionTest, BatchingWithPaddingRejectsMismatchedRanks) { + BasicBatchScheduler::Options schedule_options; + schedule_options.max_batch_size = 2; + schedule_options.batch_timeout_micros = 1e6; + schedule_options.num_batch_threads = 1; + std::unique_ptr batching_session; + BatchingSessionOptions batching_session_options; + batching_session_options.pad_variable_length_inputs = true; + TF_ASSERT_OK(CreateBasicBatchingSession( + schedule_options, batching_session_options, {{"x"}, {"y"}}, + CreateMatrixHalfPlusTwoSession(), &batching_session)); + + auto expect_rank_error = [&batching_session](Tensor input) { + std::vector outputs; + absl::Status status = + batching_session->Run({{"x", input}}, {"y"}, {}, &outputs); + EXPECT_EQ(status.code(), absl::StatusCode::kFailedPrecondition); + EXPECT_THAT(status.message(), HasSubstr("different ranks")); + }; + std::unique_ptr first_request_thread(Env::Default()->StartThread( + ThreadOptions(), "first_request", [&expect_rank_error] { + expect_rank_error(test::AsTensor({1, 2}, {1, 2})); + })); + std::unique_ptr second_request_thread(Env::Default()->StartThread( + ThreadOptions(), "second_request", [&expect_rank_error] { + expect_rank_error(test::AsTensor({3, 4}, {1, 1, 2})); + })); +} + TEST_P(BatchingSessionTest, BatchingWithLargeBatch) { BasicBatchScheduler::Options schedule_options; schedule_options.max_batch_size = 3; diff --git a/tensorflow_serving/batching/batching_util.cc b/tensorflow_serving/batching/batching_util.cc index 55af21ff1e4..9349a5c97e9 100644 --- a/tensorflow_serving/batching/batching_util.cc +++ b/tensorflow_serving/batching/batching_util.cc @@ -153,26 +153,51 @@ absl::Status PadTensorOfSpecificType(const Tensor& tensor, } } -std::map> CalculateMaxDimSizes( - const std::vector>>& batch) { - std::map> max_dim_sizes; +absl::Status CalculateMaxDimSizes( + const std::vector>>& batch, + std::map>* max_dim_sizes) { + if (batch.empty()) { + return absl::InvalidArgumentError( + "Cannot calculate dimensions for an empty batch."); + } + max_dim_sizes->clear(); // Populate 'max_dim_sizes' // init const std::vector>& task_inputs = batch[0]; for (const auto& entry : task_inputs) { const std::string& tensor_name = entry.first; const Tensor& tensor = entry.second; - max_dim_sizes[tensor_name] = std::vector(tensor.dims(), 0); + if (!max_dim_sizes->emplace(tensor_name, std::vector(tensor.dims(), 0)) + .second) { + return absl::FailedPreconditionError(absl::StrCat( + "Task has duplicate input tensor name '", tensor_name, "'.")); + } } // fill for (int i = 0; i < batch.size(); ++i) { const std::vector>& task_inputs = batch[i]; + if (task_inputs.size() != max_dim_sizes->size()) { + return absl::FailedPreconditionError( + "Tasks in a single batch have different numbers of input tensors."); + } for (const auto& entry : task_inputs) { const std::string& tensor_name = entry.first; const Tensor& tensor = entry.second; - std::vector& max_dim_sizes_for_one_input = - max_dim_sizes[tensor_name]; + auto max_dim_sizes_it = max_dim_sizes->find(tensor_name); + if (max_dim_sizes_it == max_dim_sizes->end()) { + return absl::FailedPreconditionError(absl::StrCat( + "Tasks in a single batch have different input tensor names; '", + tensor_name, "' was not present in the first task.")); + } + std::vector& max_dim_sizes_for_one_input = max_dim_sizes_it->second; + if (static_cast(tensor.dims()) != + max_dim_sizes_for_one_input.size()) { + return absl::FailedPreconditionError(absl::StrCat( + "Tensors with name '", tensor_name, + "' from different tasks have different ranks: expected ", + max_dim_sizes_for_one_input.size(), ", got ", tensor.dims(), ".")); + } for (int j = 0; j < tensor.dims(); ++j) { const int old_max_size = max_dim_sizes_for_one_input[j]; if (tensor.shape().dim_size(j) > old_max_size) { @@ -181,12 +206,17 @@ std::map> CalculateMaxDimSizes( } } } - return max_dim_sizes; + return absl::OkStatus(); } absl::Status AddPadding(const Tensor& tensor, absl::Span max_dim_sizes, Tensor* padded_tensor) { + if (static_cast(tensor.dims()) != max_dim_sizes.size()) { + return absl::InvalidArgumentError(absl::StrCat( + "Tensor rank ", tensor.dims(), + " does not match maximum-dimension rank ", max_dim_sizes.size(), ".")); + } const DataType input_dtype = tensor.dtype(); absl::Status padding_status; #define CASE(type) \ diff --git a/tensorflow_serving/batching/batching_util.h b/tensorflow_serving/batching/batching_util.h index 246d4cec6fa..f61833c2252 100644 --- a/tensorflow_serving/batching/batching_util.h +++ b/tensorflow_serving/batching/batching_util.h @@ -16,10 +16,12 @@ limitations under the License. #ifndef TENSORFLOW_SERVING_BATCHING_BATCHING_UTIL_H_ #define TENSORFLOW_SERVING_BATCHING_BATCHING_UTIL_H_ +#include #include #include #include +#include "absl/status/status.h" #include "absl/strings/str_cat.h" #include "absl/types/span.h" #include "tensorflow/core/framework/tensor.h" @@ -40,8 +42,9 @@ namespace serving { // // the following map will be generated: // {'tensor_a': [200, 500, 400], 'tensor_b': [200]} -std::map> CalculateMaxDimSizes( - const std::vector>>& batch); +absl::Status CalculateMaxDimSizes( + const std::vector>>& batch, + std::map>* max_dim_sizes); // Pads tensor so that its shape becomes as specified in max_dim_sizes, // except for zeroth dimension, which is left as is. diff --git a/tensorflow_serving/batching/batching_util_test.cc b/tensorflow_serving/batching/batching_util_test.cc index 58977a4b59a..05305927eac 100644 --- a/tensorflow_serving/batching/batching_util_test.cc +++ b/tensorflow_serving/batching/batching_util_test.cc @@ -32,6 +32,7 @@ namespace serving { namespace { using ::testing::ElementsAre; +using ::testing::HasSubstr; using ::testing::Pair; using ::testing::UnorderedElementsAre; @@ -56,13 +57,31 @@ TEST(BatchingUtilTest, CalculateMaxDimSizes) { CreateInputsWithTensorShapes(shapes2); std::vector>> batch{inputs1, inputs2}; - std::map> max_dim_sizes = - CalculateMaxDimSizes(batch); + std::map> max_dim_sizes; + ASSERT_EQ(absl::OkStatus(), + CalculateMaxDimSizes(batch, &max_dim_sizes)); EXPECT_THAT(max_dim_sizes, UnorderedElementsAre(Pair("x0", ElementsAre(20, 50, 30)), Pair("x1", ElementsAre(20, 101)))); } +TEST(BatchingUtilTest, CalculateMaxDimSizesRejectsMismatchedRanks) { + const auto rank_two = CreateInputsWithTensorShapes({TensorShape({1, 2})}); + const auto rank_three = + CreateInputsWithTensorShapes({TensorShape({1, 1, 2})}); + + for (const auto& batch : + {std::vector>>{rank_two, + rank_three}, + std::vector>>{rank_three, + rank_two}}) { + std::map> max_dim_sizes; + absl::Status status = CalculateMaxDimSizes(batch, &max_dim_sizes); + EXPECT_EQ(status.code(), absl::StatusCode::kFailedPrecondition); + EXPECT_THAT(status.message(), HasSubstr("different ranks")); + } +} + TEST(BatchingUtilTest, AddPadding) { const std::vector max_dim_sizes{20, 100, 200}; const std::vector types{ @@ -98,6 +117,18 @@ TEST(BatchingUtilTest, AddPaddingTensorWithUnsupportedRank) { "Only tensors with rank from 1 to 6 can be padded."), AddPadding(tensor, max_dim_sizes, &padded_tensor)); } + +TEST(BatchingUtilTest, AddPaddingRejectsMismatchedRank) { + Tensor padded_tensor; + absl::Status status = + AddPadding(Tensor(DT_FLOAT, {1, 2}), {1, 2, 3}, &padded_tensor); + EXPECT_EQ(status.code(), absl::StatusCode::kInvalidArgument); + EXPECT_THAT(status.message(), HasSubstr("does not match")); + + status = AddPadding(Tensor(DT_FLOAT, {1, 2, 3}), {1, 2}, &padded_tensor); + EXPECT_EQ(status.code(), absl::StatusCode::kInvalidArgument); + EXPECT_THAT(status.message(), HasSubstr("does not match")); +} } // namespace } // namespace serving } // namespace tensorflow diff --git a/tensorflow_serving/batching/tfrt_saved_model_with_batching.cc b/tensorflow_serving/batching/tfrt_saved_model_with_batching.cc index f9f04725129..90c3e91f9b0 100644 --- a/tensorflow_serving/batching/tfrt_saved_model_with_batching.cc +++ b/tensorflow_serving/batching/tfrt_saved_model_with_batching.cc @@ -237,11 +237,20 @@ absl::Status SavedModelWithBatching::Run( // TODO(b/168220822): Once tfrt supports tensor split/pad/concat utilities and // removes llvm dependency, refactors this function accordingly (return type may // change). -std::vector> CalculateMaxDimSizes( - const Batch& batch) { - std::vector> max_dim_sizes; +absl::Status CalculateMaxDimSizes( + const Batch& batch, + std::vector>* max_dim_sizes) { + if (batch.num_tasks() < 1) { + return absl::InvalidArgumentError( + "Cannot calculate dimensions for an empty batch."); + } + max_dim_sizes->clear(); for (int batch_idx = 0; batch_idx < batch.num_tasks(); ++batch_idx) { const auto inputs = batch.task(batch_idx).tfrt_inputs; + if (batch_idx > 0 && inputs.size() != max_dim_sizes->size()) { + return absl::FailedPreconditionError( + "Tasks in a single batch have different numbers of input tensors."); + } for (int tensor_idx = 0; tensor_idx < inputs.size(); ++tensor_idx) { const Tensor& tensor = inputs[tensor_idx]; const TensorShape& shape = tensor.shape(); @@ -254,16 +263,23 @@ std::vector> CalculateMaxDimSizes( } if (batch_idx == 0) { - max_dim_sizes.push_back(std::move(dims)); + max_dim_sizes->push_back(std::move(dims)); } else { + absl::InlinedVector& max_sizes = (*max_dim_sizes)[tensor_idx]; + if (max_sizes.size() != static_cast(rank)) { + return absl::FailedPreconditionError(absl::StrCat( + "Tensors at input index ", tensor_idx, + " from different tasks have different ranks: expected ", + max_sizes.size(), ", got ", rank, ".")); + } for (int rank_idx = 0; rank_idx < rank; ++rank_idx) { - int& cur_max_size = max_dim_sizes[tensor_idx][rank_idx]; + int& cur_max_size = max_sizes[rank_idx]; cur_max_size = std::max(cur_max_size, dims[rank_idx]); } } } } - return max_dim_sizes; + return absl::OkStatus(); } absl::Status SavedModelWithBatching::BatchInputTensors( @@ -282,7 +298,7 @@ absl::Status SavedModelWithBatching::BatchInputTensors( std::vector> max_dim_sizes; if (options_.pad_variable_length_inputs) { - max_dim_sizes = CalculateMaxDimSizes(batch); + TF_RETURN_IF_ERROR(CalculateMaxDimSizes(batch, &max_dim_sizes)); } // TODO(b/168220822): Padding logic below operates on tfrt inputs. It's pretty diff --git a/tensorflow_serving/batching/tfrt_saved_model_with_batching_test.cc b/tensorflow_serving/batching/tfrt_saved_model_with_batching_test.cc index b8811c7b065..528da549641 100644 --- a/tensorflow_serving/batching/tfrt_saved_model_with_batching_test.cc +++ b/tensorflow_serving/batching/tfrt_saved_model_with_batching_test.cc @@ -383,6 +383,38 @@ TEST_F(SavedModelWithBatchingTest, BatchingWithPadding) { })); } +TEST_F(SavedModelWithBatchingTest, BatchingWithPaddingRejectsMismatchedRanks) { + Initialize(BuildSchedulerOptions(/*max_batch_size=*/2), + BuildSavedModelBatchingOptions( + /*pad_variable_length_inputs=*/true, + /*allowed_batch_sizes=*/{})); + + auto inputs = MakeTensorsBatch({ + {{{1, 2}, TensorShape({1, 2})}}, + {{{3, 4}, TensorShape({1, 1, 2})}}, + }); + + EXPECT_CALL( + *wrapped_saved_model_, + Run(_, kFunctionOne, ::testing::An>(), _)) + .Times(0); + + tfrt::SavedModel::RunOptions run_options; + auto expect_rank_error = [this, &inputs, &run_options](int input_index) { + std::vector outputs; + absl::Status status = saved_model_with_batching_->Run( + run_options, kFunctionOne, inputs[input_index], &outputs); + EXPECT_THAT(status, + TFStatusIs(error::FAILED_PRECONDITION, "different ranks")); + }; + std::unique_ptr first_request_thread(Env::Default()->StartThread( + ThreadOptions(), "first_request_thread", + [&expect_rank_error] { expect_rank_error(0); })); + std::unique_ptr second_request_thread(Env::Default()->StartThread( + ThreadOptions(), "second_request_thread", + [&expect_rank_error] { expect_rank_error(1); })); +} + // Tests that batching tensors with variable length dimension size (except for // batching dimension) returns an appropriate error when padding is turned off. TEST_F(SavedModelWithBatchingTest, UnequalShapesWhenPaddingIsTurnedOff) {