summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorPisarenko Grigoriy <[email protected]>2024-02-20 20:02:26 +0300
committerGitHub <[email protected]>2024-02-20 20:02:26 +0300
commite134eb366ebd22e646fc81c1d7428d0e6736eafe (patch)
tree65bf7ed82d95fbd8fdf4359ffc680c4ba203a7d0
parentff55aa5ad5c364ba08e704643f8c943003d6e479 (diff)
YQ-2883 fix retries for recover point COMPLETING (#2107)
-rw-r--r--ydb/core/fq/libs/compute/ydb/actors_factory.cpp5
-rw-r--r--ydb/core/fq/libs/compute/ydb/actors_factory.h3
-rw-r--r--ydb/core/fq/libs/compute/ydb/result_writer_actor.cpp14
-rw-r--r--ydb/core/fq/libs/compute/ydb/result_writer_actor.h1
-rw-r--r--ydb/core/fq/libs/compute/ydb/ydb_run_actor.cpp4
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);
}