diff options
| author | Pisarenko Grigoriy <[email protected]> | 2024-02-20 20:02:26 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-02-20 20:02:26 +0300 |
| commit | e134eb366ebd22e646fc81c1d7428d0e6736eafe (patch) | |
| tree | 65bf7ed82d95fbd8fdf4359ffc680c4ba203a7d0 | |
| parent | ff55aa5ad5c364ba08e704643f8c943003d6e479 (diff) | |
YQ-2883 fix retries for recover point COMPLETING (#2107)
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/actors_factory.cpp | 5 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/actors_factory.h | 3 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp | 14 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/result_writer_actor.h | 1 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp | 4 |
5 files changed, 20 insertions, 7 deletions
diff --git a/ydb/core/fq/libs/compute/ydb/actors_factory.cpp b/ydb/core/fq/libs/compute/ydb/actors_factory.cpp index 9db333a97da..fc6d2606c0f 100644 --- a/ydb/core/fq/libs/compute/ydb/actors_factory.cpp +++ b/ydb/core/fq/libs/compute/ydb/actors_factory.cpp @@ -59,8 +59,9 @@ struct TActorFactory : public IActorFactory { std::unique_ptr<NActors::IActor> CreateResultWriter(const NActors::TActorId& parent, const NActors::TActorId& connector, const NActors::TActorId& pinger, - const NKikimr::NOperationId::TOperationId& operationId) const override { - return CreateResultWriterActor(Params, parent, connector, pinger, operationId, Counters); + const NKikimr::NOperationId::TOperationId& operationId, + bool operationEntryExpected) const override { + return CreateResultWriterActor(Params, parent, connector, pinger, operationId, operationEntryExpected, Counters); } std::unique_ptr<NActors::IActor> CreateResourcesCleaner(const NActors::TActorId& parent, diff --git a/ydb/core/fq/libs/compute/ydb/actors_factory.h b/ydb/core/fq/libs/compute/ydb/actors_factory.h index ae85da060f7..4c91bfc1b6b 100644 --- a/ydb/core/fq/libs/compute/ydb/actors_factory.h +++ b/ydb/core/fq/libs/compute/ydb/actors_factory.h @@ -28,7 +28,8 @@ struct IActorFactory : public TThrRefBase { virtual std::unique_ptr<NActors::IActor> CreateResultWriter(const NActors::TActorId& parent, const NActors::TActorId& connector, const NActors::TActorId& pinger, - const NKikimr::NOperationId::TOperationId& operationId) const = 0; + const NKikimr::NOperationId::TOperationId& operationId, + bool operationEntryExpected) const = 0; virtual std::unique_ptr<NActors::IActor> CreateResourcesCleaner(const NActors::TActorId& parent, const NActors::TActorId& connector, const NYdb::TOperation::TOperationId& operationId) const = 0; diff --git a/ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp b/ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp index ef6da5653b4..b6f0ff8efc0 100644 --- a/ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp @@ -202,13 +202,14 @@ public: } }; - TResultWriterActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NKikimr::NOperationId::TOperationId& operationId, const ::NYql::NCommon::TServiceCounters& queryCounters) + TResultWriterActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NKikimr::NOperationId::TOperationId& operationId, bool operationEntryExpected, const ::NYql::NCommon::TServiceCounters& queryCounters) : TBaseComputeActor(queryCounters, "ResultWriter") , Params(params) , Parent(parent) , Connector(connector) , Pinger(pinger) , OperationId(operationId) + , OperationEntryExpected(operationEntryExpected) , Counters(GetStepCountersSubgroup()) {} @@ -246,6 +247,13 @@ public: void Handle(const TEvYdbCompute::TEvGetOperationResponse::TPtr& ev) { const auto& response = *ev.Get()->Get(); + if (!OperationEntryExpected && response.Status == NYdb::EStatus::NOT_FOUND) { + LOG_I("Operation has been already removed"); + Send(Parent, new TEvYdbCompute::TEvResultWriterResponse({}, NYdb::EStatus::SUCCESS)); + CompleteAndPassAway(); + return; + } + if (response.Status != NYdb::EStatus::SUCCESS) { LOG_E("Can't get operation: " << ev->Get()->Issues.ToOneLineString()); Send(Parent, new TEvYdbCompute::TEvResultWriterResponse(ev->Get()->Issues, ev->Get()->Status)); @@ -314,6 +322,7 @@ private: TActorId Connector; TActorId Pinger; NKikimr::NOperationId::TOperationId OperationId; + const bool OperationEntryExpected; TCounters Counters; TInstant StartTime; TString FetchToken; @@ -325,8 +334,9 @@ std::unique_ptr<NActors::IActor> CreateResultWriterActor(const TRunActorParams& const TActorId& connector, const TActorId& pinger, const NKikimr::NOperationId::TOperationId& operationId, + bool operationEntryExpected, const ::NYql::NCommon::TServiceCounters& queryCounters) { - return std::make_unique<TResultWriterActor>(params, parent, connector, pinger, operationId, queryCounters); + return std::make_unique<TResultWriterActor>(params, parent, connector, pinger, operationId, operationEntryExpected, queryCounters); } } diff --git a/ydb/core/fq/libs/compute/ydb/result_writer_actor.h b/ydb/core/fq/libs/compute/ydb/result_writer_actor.h index ee24d14772b..ca6c1454d42 100644 --- a/ydb/core/fq/libs/compute/ydb/result_writer_actor.h +++ b/ydb/core/fq/libs/compute/ydb/result_writer_actor.h @@ -13,6 +13,7 @@ std::unique_ptr<NActors::IActor> CreateResultWriterActor(const TRunActorParams& const NActors::TActorId& connector, const NActors::TActorId& pinger, const NKikimr::NOperationId::TOperationId& operationId, + bool operationEntryExpected, const ::NYql::NCommon::TServiceCounters& queryCounters); } diff --git a/ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp b/ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp index 41701816a81..a18d2043c5b 100644 --- a/ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp @@ -118,7 +118,7 @@ public: Params.Status = response.ComputeStatus; LOG_I("StatusTrackerResponse (success) " << response.Status << " ExecStatus: " << static_cast<int>(response.ExecStatus) << " Issues: " << response.Issues.ToOneLineString()); if (response.ExecStatus == NYdb::NQuery::EExecStatus::Completed) { - Register(ActorFactory->CreateResultWriter(SelfId(), Connector, Pinger, Params.OperationId).release()); + Register(ActorFactory->CreateResultWriter(SelfId(), Connector, Pinger, Params.OperationId, true).release()); } else { CreateResourcesCleaner(); } @@ -192,7 +192,7 @@ public: break; case FederatedQuery::QueryMeta::COMPLETING: if (Params.OperationId.GetKind() != Ydb::TOperationId::UNUSED) { - Register(ActorFactory->CreateResultWriter(SelfId(), Connector, Pinger, Params.OperationId).release()); + Register(ActorFactory->CreateResultWriter(SelfId(), Connector, Pinger, Params.OperationId, false).release()); } else { CreateFinalizer(Params.Status); } |
