diff --git a/src/ray/raylet/local_object_manager.cc b/src/ray/raylet/local_object_manager.cc index 55d57c8f4..721adb6bd 100644 --- a/src/ray/raylet/local_object_manager.cc +++ b/src/ray/raylet/local_object_manager.cc @@ -306,6 +306,8 @@ void LocalObjectManager::AsyncRestoreSpilledObject( // If the same object is restoring, we dedup here. return; } + RAY_CHECK(objects_pending_restore_.emplace(object_id).second) + << "Object dedupe wasn't done properly. Please report if you see this issue."; io_worker_pool_.PopRestoreWorker([this, object_id, object_url, callback]( std::shared_ptr io_worker) { auto start_time = absl::GetCurrentTimeNanos(); @@ -313,8 +315,6 @@ void LocalObjectManager::AsyncRestoreSpilledObject( rpc::RestoreSpilledObjectsRequest request; request.add_spilled_objects_url(std::move(object_url)); request.add_object_ids_to_restore(object_id.Binary()); - RAY_CHECK(objects_pending_restore_.emplace(object_id).second) - << "Object dedupe wasn't done properly. Please report if you see this issue."; io_worker->rpc_client()->RestoreSpilledObjects( request, [this, start_time, object_id, callback, io_worker]( diff --git a/src/ray/raylet/test/local_object_manager_test.cc b/src/ray/raylet/test/local_object_manager_test.cc index 6a5ec3a45..616e73482 100644 --- a/src/ray/raylet/test/local_object_manager_test.cc +++ b/src/ray/raylet/test/local_object_manager_test.cc @@ -150,7 +150,7 @@ class MockIOWorkerPool : public IOWorkerPoolInterface { void PopRestoreWorker( std::function)> callback) override { - callback(io_worker); + restoration_callbacks.push_back(callback); } void PopDeleteWorker( @@ -158,6 +158,17 @@ class MockIOWorkerPool : public IOWorkerPoolInterface { callback(io_worker); } + bool RestoreWorkerPushed() { + if (restoration_callbacks.size() == 0) { + return false; + } + const auto callback = restoration_callbacks.front(); + callback(io_worker); + restoration_callbacks.pop_front(); + return true; + } + + std::list)>> restoration_callbacks; std::shared_ptr io_worker_client = std::make_shared(); std::shared_ptr io_worker = @@ -322,6 +333,13 @@ TEST_F(LocalObjectManagerTest, TestRestoreSpilledObject) { num_times_fired++; }); } + ASSERT_EQ(num_times_fired, 0); + + // When restore workers are pushed, the request should be dedupped. + for (int i = 0; i < 10; i++) { + worker_pool.RestoreWorkerPushed(); + ASSERT_EQ(num_times_fired, 0); + } worker_pool.io_worker_client->ReplyRestoreObjects(10); ASSERT_EQ(num_times_fired, 1); }