From 1a8d0af814a130f327a33eb41dd57defdb60589b Mon Sep 17 00:00:00 2001 From: Stephanie Wang Date: Wed, 26 Jun 2019 11:21:00 -0700 Subject: [PATCH] Remove debug check for uncommitted lineage (#5038) --- src/ray/raylet/lineage_cache.cc | 9 +++---- src/ray/raylet/lineage_cache.h | 7 +++--- src/ray/raylet/lineage_cache_test.cc | 36 ++++++++++++++-------------- src/ray/raylet/node_manager.cc | 15 ++++-------- 4 files changed, 31 insertions(+), 36 deletions(-) diff --git a/src/ray/raylet/lineage_cache.cc b/src/ray/raylet/lineage_cache.cc index 68d5aa817..9d700b805 100644 --- a/src/ray/raylet/lineage_cache.cc +++ b/src/ray/raylet/lineage_cache.cc @@ -245,8 +245,8 @@ void GetUncommittedLineageHelper(const TaskID &task_id, const Lineage &lineage_f } } -Lineage LineageCache::GetUncommittedLineageOrDie(const TaskID &task_id, - const ClientID &node_id) const { +Lineage LineageCache::GetUncommittedLineage(const TaskID &task_id, + const ClientID &node_id) const { Lineage uncommitted_lineage; // Add all uncommitted ancestors from the lineage cache to the uncommitted // lineage of the requested task. @@ -256,8 +256,9 @@ Lineage LineageCache::GetUncommittedLineageOrDie(const TaskID &task_id, // already explicitly forwarded to this node before. if (uncommitted_lineage.GetEntries().empty()) { auto entry = lineage_.GetEntry(task_id); - RAY_CHECK(entry); - RAY_CHECK(uncommitted_lineage.SetEntry(entry->TaskData(), entry->GetStatus())); + if (entry) { + RAY_CHECK(uncommitted_lineage.SetEntry(entry->TaskData(), entry->GetStatus())); + } } return uncommitted_lineage; } diff --git a/src/ray/raylet/lineage_cache.h b/src/ray/raylet/lineage_cache.h index 37ce5caf6..6ff3fb5fc 100644 --- a/src/ray/raylet/lineage_cache.h +++ b/src/ray/raylet/lineage_cache.h @@ -244,13 +244,12 @@ class LineageCache { /// The uncommitted lineage consists of all tasks in the given task's lineage /// that have not been committed in the GCS, as far as we know. /// - /// \param task_id The ID of the task to get the uncommitted lineage for. It is - /// a fatal error if the task is not found. + /// \param task_id The ID of the task to get the uncommitted lineage for. If + /// the task is not found, then the returned lineage will be empty. /// \param node_id The ID of the receiving node. /// \return The uncommitted, unforwarded lineage of the task. The returned lineage /// includes the entry for the requested entry_id. - Lineage GetUncommittedLineageOrDie(const TaskID &task_id, - const ClientID &node_id) const; + Lineage GetUncommittedLineage(const TaskID &task_id, const ClientID &node_id) const; /// Handle the commit of a task entry in the GCS. This attempts to evict the /// task if possible. diff --git a/src/ray/raylet/lineage_cache_test.cc b/src/ray/raylet/lineage_cache_test.cc index a6184902f..27a231a85 100644 --- a/src/ray/raylet/lineage_cache_test.cc +++ b/src/ray/raylet/lineage_cache_test.cc @@ -167,7 +167,7 @@ std::vector InsertTaskChain(LineageCache &lineage_cache, return arguments; } -TEST_F(LineageCacheTest, TestGetUncommittedLineageOrDie) { +TEST_F(LineageCacheTest, TestGetUncommittedLineage) { // Insert two independent chains of tasks. std::vector tasks1; auto return_values1 = @@ -187,7 +187,7 @@ TEST_F(LineageCacheTest, TestGetUncommittedLineageOrDie) { // Get the uncommitted lineage for the last task (the leaf) of one of the chains. auto uncommitted_lineage = - lineage_cache_.GetUncommittedLineageOrDie(task_ids1.back(), ClientID::Nil()); + lineage_cache_.GetUncommittedLineage(task_ids1.back(), ClientID::Nil()); // Check that the uncommitted lineage is exactly equal to the first chain of tasks. ASSERT_EQ(task_ids1.size(), uncommitted_lineage.GetEntries().size()); for (auto &task_id : task_ids1) { @@ -207,8 +207,8 @@ TEST_F(LineageCacheTest, TestGetUncommittedLineageOrDie) { } // Get the uncommitted lineage for the inserted task. - uncommitted_lineage = lineage_cache_.GetUncommittedLineageOrDie( - combined_task_ids.back(), ClientID::Nil()); + uncommitted_lineage = + lineage_cache_.GetUncommittedLineage(combined_task_ids.back(), ClientID::Nil()); // Check that the uncommitted lineage is exactly equal to the entire set of // tasks inserted so far. ASSERT_EQ(combined_task_ids.size(), uncommitted_lineage.GetEntries().size()); @@ -262,9 +262,9 @@ TEST_F(LineageCacheTest, TestMarkTaskAsForwarded) { lineage_cache_.MarkTaskAsForwarded(forwarded_task_id, node_id); auto uncommitted_lineage = - lineage_cache_.GetUncommittedLineageOrDie(remaining_task_id, node_id); + lineage_cache_.GetUncommittedLineage(remaining_task_id, node_id); auto uncommitted_lineage_all = - lineage_cache_.GetUncommittedLineageOrDie(remaining_task_id, node_id2); + lineage_cache_.GetUncommittedLineage(remaining_task_id, node_id2); ASSERT_EQ(1, uncommitted_lineage.GetEntries().size()); ASSERT_EQ(4, uncommitted_lineage_all.GetEntries().size()); @@ -273,7 +273,7 @@ TEST_F(LineageCacheTest, TestMarkTaskAsForwarded) { // Check that lineage of requested task includes itself, regardless of whether // it has been forwarded before. auto uncommitted_lineage_forwarded = - lineage_cache_.GetUncommittedLineageOrDie(forwarded_task_id, node_id); + lineage_cache_.GetUncommittedLineage(forwarded_task_id, node_id); ASSERT_EQ(1, uncommitted_lineage_forwarded.GetEntries().size()); } @@ -333,8 +333,8 @@ TEST_F(LineageCacheTest, TestEvictChain) { // the flushed task, but its lineage should not be evicted yet. mock_gcs_.Flush(); ASSERT_EQ(lineage_cache_ - .GetUncommittedLineageOrDie(tasks.back().GetTaskSpecification().TaskId(), - ClientID::Nil()) + .GetUncommittedLineage(tasks.back().GetTaskSpecification().TaskId(), + ClientID::Nil()) .GetEntries() .size(), tasks.size()); @@ -346,8 +346,8 @@ TEST_F(LineageCacheTest, TestEvictChain) { mock_gcs_.RemoteAdd(tasks.at(1).GetTaskSpecification().TaskId(), task_data)); mock_gcs_.Flush(); ASSERT_EQ(lineage_cache_ - .GetUncommittedLineageOrDie(tasks.back().GetTaskSpecification().TaskId(), - ClientID::Nil()) + .GetUncommittedLineage(tasks.back().GetTaskSpecification().TaskId(), + ClientID::Nil()) .GetEntries() .size(), tasks.size()); @@ -386,8 +386,8 @@ TEST_F(LineageCacheTest, TestEvictManyParents) { mock_gcs_.Flush(); ASSERT_EQ(lineage_cache_.GetLineage().GetEntries().size(), total_tasks); ASSERT_EQ(lineage_cache_ - .GetUncommittedLineageOrDie(child_task.GetTaskSpecification().TaskId(), - ClientID::Nil()) + .GetUncommittedLineage(child_task.GetTaskSpecification().TaskId(), + ClientID::Nil()) .GetEntries() .size(), total_tasks); @@ -402,8 +402,8 @@ TEST_F(LineageCacheTest, TestEvictManyParents) { // since the parent tasks have no dependencies. ASSERT_EQ(lineage_cache_.GetLineage().GetEntries().size(), total_tasks); ASSERT_EQ(lineage_cache_ - .GetUncommittedLineageOrDie( - child_task.GetTaskSpecification().TaskId(), ClientID::Nil()) + .GetUncommittedLineage(child_task.GetTaskSpecification().TaskId(), + ClientID::Nil()) .GetEntries() .size(), total_tasks); @@ -427,7 +427,7 @@ TEST_F(LineageCacheTest, TestEviction) { // uncommitted lineage. const auto last_task_id = tasks.back().GetTaskSpecification().TaskId(); auto uncommitted_lineage = - lineage_cache_.GetUncommittedLineageOrDie(last_task_id, ClientID::Nil()); + lineage_cache_.GetUncommittedLineage(last_task_id, ClientID::Nil()); ASSERT_EQ(uncommitted_lineage.GetEntries().size(), lineage_size); // Simulate executing the first task on a remote node and adding it to the @@ -461,7 +461,7 @@ TEST_F(LineageCacheTest, TestEviction) { // All tasks have now been flushed. Check that enough lineage has been // evicted that the uncommitted lineage is now less than the maximum size. uncommitted_lineage = - lineage_cache_.GetUncommittedLineageOrDie(last_task_id, ClientID::Nil()); + lineage_cache_.GetUncommittedLineage(last_task_id, ClientID::Nil()); ASSERT_TRUE(uncommitted_lineage.GetEntries().size() < max_lineage_size_); // The remaining task should have no uncommitted lineage. ASSERT_EQ(uncommitted_lineage.GetEntries().size(), 1); @@ -481,7 +481,7 @@ TEST_F(LineageCacheTest, TestOutOfOrderEviction) { // uncommitted lineage. const auto last_task_id = tasks.back().GetTaskSpecification().TaskId(); auto uncommitted_lineage = - lineage_cache_.GetUncommittedLineageOrDie(last_task_id, ClientID::Nil()); + lineage_cache_.GetUncommittedLineage(last_task_id, ClientID::Nil()); ASSERT_EQ(uncommitted_lineage.GetEntries().size(), lineage_size); ASSERT_EQ(lineage_cache_.GetLineage().GetEntries().size(), lineage_size); diff --git a/src/ray/raylet/node_manager.cc b/src/ray/raylet/node_manager.cc index 226a8fb6d..106c83832 100644 --- a/src/ray/raylet/node_manager.cc +++ b/src/ray/raylet/node_manager.cc @@ -2217,22 +2217,17 @@ void NodeManager::ForwardTask( return; } - // Get and serialize the task's unforwarded, uncommitted lineage. - Lineage uncommitted_lineage; - if (lineage_cache_.ContainsTask(task_id)) { - uncommitted_lineage = lineage_cache_.GetUncommittedLineageOrDie(task_id, node_id); - } else { - // TODO: We expected the lineage to be in cache, but it was evicted (#3813). - // This is a bug but is not fatal to the application. - RAY_DCHECK(false) << "No lineage cache entry found for task " << task_id; + // Get the task's unforwarded, uncommitted lineage. + Lineage uncommitted_lineage = lineage_cache_.GetUncommittedLineage(task_id, node_id); + if (uncommitted_lineage.GetEntries().empty()) { + // There is no uncommitted lineage. This can happen if the lineage was + // already evicted before we forwarded the task. uncommitted_lineage.SetEntry(task, GcsStatus::NONE); } auto entry = uncommitted_lineage.GetEntryMutable(task_id); Task &lineage_cache_entry_task = entry->TaskDataMutable(); - // Increment forward count for the forwarded task. lineage_cache_entry_task.IncrementNumForwards(); - RAY_LOG(DEBUG) << "Forwarding task " << task_id << " from " << gcs_client_->client_table().GetLocalClientId() << " to " << node_id << " spillback="