summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorPisarenko Grigoriy <[email protected]>2024-02-07 00:21:44 +0300
committerGitHub <[email protected]>2024-02-07 00:21:44 +0300
commitd010ed3f89cde39cb409fbe2448e871bced2eb43 (patch)
tree17934e9d2122a9defaa4ef0c842e648a918cd847
parente889ee1e6fc13f3ea314b24cf412f371fb09b9da (diff)
YQ-2734 added retries for internal queries (#1457)
* Added retries for internal queries * Moved logic to query_actor.h
-rw-r--r--ydb/core/kqp/common/events/script_executions.h89
-rw-r--r--ydb/core/kqp/executer_actor/kqp_data_executer.cpp4
-rw-r--r--ydb/core/kqp/finalize_script_service/kqp_finalize_script_actor.cpp22
-rw-r--r--ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.cpp7
-rw-r--r--ydb/core/kqp/proxy_service/kqp_script_executions.cpp223
-rw-r--r--ydb/core/kqp/proxy_service/kqp_script_executions.h4
-rw-r--r--ydb/library/query_actor/query_actor.h100
7 files changed, 236 insertions, 213 deletions
diff --git a/ydb/core/kqp/common/events/script_executions.h b/ydb/core/kqp/common/events/script_executions.h
index 6b2b331e368..f5157a1a10b 100644
--- a/ydb/core/kqp/common/events/script_executions.h
+++ b/ydb/core/kqp/common/events/script_executions.h
@@ -221,20 +221,28 @@ struct TEvFetchScriptResultsQueryResponse : public NActors::TEventLocal<TEvFetch
};
struct TEvSaveScriptExternalEffectRequest : public NActors::TEventLocal<TEvSaveScriptExternalEffectRequest, TKqpScriptExecutionEvents::EvSaveScriptExternalEffectRequest> {
+ struct TDescription {
+ TDescription(const TString& executionId, const TString& database, const TString& customerSuppliedId, const TString& userToken)
+ : ExecutionId(executionId)
+ , Database(database)
+ , CustomerSuppliedId(customerSuppliedId)
+ , UserToken(userToken)
+ {}
+
+ TString ExecutionId;
+ TString Database;
+
+ TString CustomerSuppliedId;
+ TString UserToken;
+ std::vector<NKqpProto::TKqpExternalSink> Sinks;
+ std::vector<TString> SecretNames;
+ };
+
TEvSaveScriptExternalEffectRequest(const TString& executionId, const TString& database, const TString& customerSuppliedId, const TString& userToken)
- : ExecutionId(executionId)
- , Database(database)
- , CustomerSuppliedId(customerSuppliedId)
- , UserToken(userToken)
+ : Description(executionId, database, customerSuppliedId, userToken)
{}
- TString ExecutionId;
- TString Database;
-
- TString CustomerSuppliedId;
- TString UserToken;
- std::vector<NKqpProto::TKqpExternalSink> Sinks;
- std::vector<TString> SecretNames;
+ TDescription Description;
};
struct TEvSaveScriptExternalEffectResponse : public NActors::TEventLocal<TEvSaveScriptExternalEffectResponse, TKqpScriptExecutionEvents::EvSaveScriptExternalEffectResponse> {
@@ -248,31 +256,41 @@ struct TEvSaveScriptExternalEffectResponse : public NActors::TEventLocal<TEvSave
};
struct TEvScriptFinalizeRequest : public NActors::TEventLocal<TEvScriptFinalizeRequest, TKqpScriptExecutionEvents::EvScriptFinalizeRequest> {
+ struct TDescription {
+ TDescription(EFinalizationStatus finalizationStatus, const TString& executionId, const TString& database,
+ Ydb::StatusIds::StatusCode operationStatus, Ydb::Query::ExecStatus execStatus, NYql::TIssues issues, std::optional<NKqpProto::TKqpStatsQuery> queryStats,
+ std::optional<TString> queryPlan, std::optional<TString> queryAst, std::optional<ui64> leaseGeneration)
+ : FinalizationStatus(finalizationStatus)
+ , ExecutionId(executionId)
+ , Database(database)
+ , OperationStatus(operationStatus)
+ , ExecStatus(execStatus)
+ , Issues(std::move(issues))
+ , QueryStats(std::move(queryStats))
+ , QueryPlan(std::move(queryPlan))
+ , QueryAst(std::move(queryAst))
+ , LeaseGeneration(leaseGeneration)
+ {}
+
+ EFinalizationStatus FinalizationStatus;
+ TString ExecutionId;
+ TString Database;
+ Ydb::StatusIds::StatusCode OperationStatus;
+ Ydb::Query::ExecStatus ExecStatus;
+ NYql::TIssues Issues;
+ std::optional<NKqpProto::TKqpStatsQuery> QueryStats;
+ std::optional<TString> QueryPlan;
+ std::optional<TString> QueryAst;
+ std::optional<ui64> LeaseGeneration;
+ };
+
TEvScriptFinalizeRequest(EFinalizationStatus finalizationStatus, const TString& executionId, const TString& database,
Ydb::StatusIds::StatusCode operationStatus, Ydb::Query::ExecStatus execStatus, NYql::TIssues issues = {}, std::optional<NKqpProto::TKqpStatsQuery> queryStats = std::nullopt,
std::optional<TString> queryPlan = std::nullopt, std::optional<TString> queryAst = std::nullopt, std::optional<ui64> leaseGeneration = std::nullopt)
- : FinalizationStatus(finalizationStatus)
- , ExecutionId(executionId)
- , Database(database)
- , OperationStatus(operationStatus)
- , ExecStatus(execStatus)
- , Issues(std::move(issues))
- , QueryStats(std::move(queryStats))
- , QueryPlan(std::move(queryPlan))
- , QueryAst(std::move(queryAst))
- , LeaseGeneration(leaseGeneration)
+ : Description(finalizationStatus, executionId, database, operationStatus, execStatus, issues, queryStats, queryPlan, queryAst, leaseGeneration)
{}
- EFinalizationStatus FinalizationStatus;
- TString ExecutionId;
- TString Database;
- Ydb::StatusIds::StatusCode OperationStatus;
- Ydb::Query::ExecStatus ExecStatus;
- NYql::TIssues Issues;
- std::optional<NKqpProto::TKqpStatsQuery> QueryStats;
- std::optional<TString> QueryPlan;
- std::optional<TString> QueryAst;
- std::optional<ui64> LeaseGeneration;
+ TDescription Description;
};
struct TEvScriptFinalizeResponse : public NActors::TEventLocal<TEvScriptFinalizeResponse, TKqpScriptExecutionEvents::EvScriptFinalizeResponse> {
@@ -284,15 +302,14 @@ struct TEvScriptFinalizeResponse : public NActors::TEventLocal<TEvScriptFinalize
};
struct TEvSaveScriptFinalStatusResponse : public NActors::TEventLocal<TEvSaveScriptFinalStatusResponse, TKqpScriptExecutionEvents::EvSaveScriptFinalStatusResponse> {
- TEvSaveScriptFinalStatusResponse(const TString& customerSuppliedId, const TString& userToken)
- : CustomerSuppliedId(customerSuppliedId)
- , UserToken(userToken)
- {}
-
+ bool ApplicateScriptExternalEffectRequired = false;
+ bool OperationAlreadyFinalized = false;
TString CustomerSuppliedId;
TString UserToken;
std::vector<NKqpProto::TKqpExternalSink> Sinks;
std::vector<TString> SecretNames;
+ Ydb::StatusIds::StatusCode Status;
+ NYql::TIssues Issues;
};
struct TEvDescribeSecretsResponse : public NActors::TEventLocal<TEvDescribeSecretsResponse, TKqpScriptExecutionEvents::EvDescribeSecretsResponse> {
diff --git a/ydb/core/kqp/executer_actor/kqp_data_executer.cpp b/ydb/core/kqp/executer_actor/kqp_data_executer.cpp
index a6fdcac9a80..0c7cdb69961 100644
--- a/ydb/core/kqp/executer_actor/kqp_data_executer.cpp
+++ b/ydb/core/kqp/executer_actor/kqp_data_executer.cpp
@@ -1595,13 +1595,13 @@ private:
for (const auto& sink : stage.GetSinks()) {
if (sink.GetTypeCase() == NKqpProto::TKqpSink::kExternalSink) {
SaveScriptExternalEffectRequired = true;
- scriptExternalEffect->Sinks.push_back(sink.GetExternalSink());
+ scriptExternalEffect->Description.Sinks.push_back(sink.GetExternalSink());
}
}
}
}
}
- scriptExternalEffect->SecretNames = SecretNames;
+ scriptExternalEffect->Description.SecretNames = SecretNames;
if (!WaitRequired()) {
return Execute();
diff --git a/ydb/core/kqp/finalize_script_service/kqp_finalize_script_actor.cpp b/ydb/core/kqp/finalize_script_service/kqp_finalize_script_actor.cpp
index 03d2a475bbe..6ffc4b58b3d 100644
--- a/ydb/core/kqp/finalize_script_service/kqp_finalize_script_actor.cpp
+++ b/ydb/core/kqp/finalize_script_service/kqp_finalize_script_actor.cpp
@@ -22,9 +22,9 @@ public:
const NKikimrConfig::TMetadataProviderConfig& metadataProviderConfig,
const std::optional<TKqpFederatedQuerySetup>& federatedQuerySetup)
: ReplyActor_(request->Sender)
- , ExecutionId_(request->Get()->ExecutionId)
- , Database_(request->Get()->Database)
- , FinalizationStatus_(request->Get()->FinalizationStatus)
+ , ExecutionId_(request->Get()->Description.ExecutionId)
+ , Database_(request->Get()->Description.Database)
+ , FinalizationStatus_(request->Get()->Description.FinalizationStatus)
, Request_(std::move(request))
, FinalizationTimeout_(TDuration::Seconds(finalizeScriptServiceConfig.GetScriptFinalizationTimeoutSeconds()))
, MaximalSecretsSnapshotWaitTime_(2 * TDuration::Seconds(metadataProviderConfig.GetRefreshPeriodSeconds()))
@@ -32,16 +32,20 @@ public:
{}
void Bootstrap() {
- Register(CreateSaveScriptFinalStatusActor(std::move(Request_)));
+ Register(CreateSaveScriptFinalStatusActor(SelfId(), std::move(Request_)));
Become(&TScriptFinalizerActor::FetchState);
}
STRICT_STFUNC(FetchState,
hFunc(TEvSaveScriptFinalStatusResponse, Handle);
- hFunc(TEvScriptExecutionFinished, Handle);
)
void Handle(TEvSaveScriptFinalStatusResponse::TPtr& ev) {
+ if (!ev->Get()->ApplicateScriptExternalEffectRequired || ev->Get()->Status != Ydb::StatusIds::SUCCESS) {
+ Reply(ev->Get()->OperationAlreadyFinalized, ev->Get()->Status, std::move(ev->Get()->Issues));
+ return;
+ }
+
Schedule(FinalizationTimeout_, new TEvents::TEvWakeup());
Become(&TScriptFinalizerActor::PrepareState);
@@ -168,7 +172,7 @@ private:
)
void FinishScriptFinalization(std::optional<Ydb::StatusIds::StatusCode> status, NYql::TIssues issues) {
- Register(CreateScriptFinalizationFinisherActor(ExecutionId_, Database_, status, std::move(issues)));
+ Register(CreateScriptFinalizationFinisherActor(SelfId(), ExecutionId_, Database_, status, std::move(issues)));
Become(&TScriptFinalizerActor::FinishState);
}
@@ -181,7 +185,11 @@ private:
}
void Handle(TEvScriptExecutionFinished::TPtr& ev) {
- Send(ReplyActor_, ev->Release());
+ Reply(ev->Get()->OperationAlreadyFinalized, ev->Get()->Status, std::move(ev->Get()->Issues));
+ }
+
+ void Reply(bool operationAlreadyFinalized, Ydb::StatusIds::StatusCode status, NYql::TIssues&& issues) {
+ Send(ReplyActor_, new TEvScriptExecutionFinished(operationAlreadyFinalized, status, std::move(issues)));
Send(MakeKqpFinalizeScriptServiceId(SelfId().NodeId()), new TEvScriptFinalizeResponse(ExecutionId_));
PassAway();
diff --git a/ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.cpp b/ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.cpp
index cf6c66d5b59..6c0868e4c42 100644
--- a/ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.cpp
+++ b/ydb/core/kqp/finalize_script_service/kqp_finalize_script_service.cpp
@@ -27,9 +27,10 @@ public:
}
void Handle(TEvSaveScriptExternalEffectRequest::TPtr& ev) {
- ev->Get()->Sinks = FilterExternalSinks(ev->Get()->Sinks);
+ auto& description = ev->Get()->Description;
+ description.Sinks = FilterExternalSinks(description.Sinks);
- if (!ev->Get()->Sinks.empty()) {
+ if (!description.Sinks.empty()) {
Register(CreateSaveScriptExternalEffectActor(std::move(ev)));
} else {
Send(ev->Sender, new TEvSaveScriptExternalEffectResponse(Ydb::StatusIds::SUCCESS, {}));
@@ -37,7 +38,7 @@ public:
}
void Handle(TEvScriptFinalizeRequest::TPtr& ev) {
- TString executionId = ev->Get()->ExecutionId;
+ TString executionId = ev->Get()->Description.ExecutionId;
if (!FinalizationRequestsQueue_.contains(executionId)) {
WaitingFinalizationExecutions_.push(executionId);
diff --git a/ydb/core/kqp/proxy_service/kqp_script_executions.cpp b/ydb/core/kqp/proxy_service/kqp_script_executions.cpp
index 56a9dae3b58..30491e1bb18 100644
--- a/ydb/core/kqp/proxy_service/kqp_script_executions.cpp
+++ b/ydb/core/kqp/proxy_service/kqp_script_executions.cpp
@@ -478,8 +478,6 @@ private:
class TScriptLeaseUpdateActor : public TActorBootstrapped<TScriptLeaseUpdateActor> {
public:
- using IRetryPolicy = IRetryPolicy<const Ydb::StatusIds::StatusCode&>;
-
TScriptLeaseUpdateActor(const TActorId& runScriptActorId, const TString& database, const TString& executionId, TDuration leaseDuration, TIntrusivePtr<TKqpCounters> counters)
: RunScriptActorId(runScriptActorId)
, Database(database)
@@ -489,44 +487,16 @@ public:
, LeaseUpdateStartTime(TInstant::Now())
{}
- void CreateScriptLeaseUpdater() {
- Register(new TScriptLeaseUpdater(Database, ExecutionId, LeaseDuration));
- }
-
void Bootstrap() {
- CreateScriptLeaseUpdater();
+ Register(new TQueryRetryActor<TScriptLeaseUpdater, TEvScriptLeaseUpdateResponse, TString, TString, TDuration>(SelfId(), Database, ExecutionId, LeaseDuration, LeaseDuration / 2));
Become(&TScriptLeaseUpdateActor::StateFunc);
}
STRICT_STFUNC(StateFunc,
hFunc(TEvScriptLeaseUpdateResponse, Handle);
- hFunc(NActors::TEvents::TEvWakeup, Wakeup);
)
- void Wakeup(NActors::TEvents::TEvWakeup::TPtr&) {
- CreateScriptLeaseUpdater();
- }
-
void Handle(TEvScriptLeaseUpdateResponse::TPtr& ev) {
- auto queryStatus = ev->Get()->Status;
- if (!ev->Get()->ExecutionEntryExists && queryStatus == Ydb::StatusIds::BAD_REQUEST || queryStatus == Ydb::StatusIds::SUCCESS) {
- Reply(std::move(ev));
- return;
- }
-
- if (RetryState == nullptr) {
- CreateRetryState();
- }
-
- const TMaybe<TDuration> delay = RetryState->GetNextRetryDelay(queryStatus);
- if (delay) {
- Schedule(*delay, new NActors::TEvents::TEvWakeup());
- } else {
- Reply(std::move(ev));
- }
- }
-
- void Reply(TEvScriptLeaseUpdateResponse::TPtr&& ev) {
if (Counters) {
Counters->ReportLeaseUpdateLatency(TInstant::Now() - LeaseUpdateStartTime);
}
@@ -534,33 +504,6 @@ public:
PassAway();
}
- static ERetryErrorClass Retryable(const Ydb::StatusIds::StatusCode& status) {
- if (status == Ydb::StatusIds::SUCCESS) {
- return ERetryErrorClass::NoRetry;
- }
-
- if (status == Ydb::StatusIds::INTERNAL_ERROR
- || status == Ydb::StatusIds::UNAVAILABLE
- || status == Ydb::StatusIds::TIMEOUT
- || status == Ydb::StatusIds::BAD_SESSION
- || status == Ydb::StatusIds::SESSION_EXPIRED
- || status == Ydb::StatusIds::SESSION_BUSY
- || status == Ydb::StatusIds::ABORTED) {
- return ERetryErrorClass::ShortRetry;
- }
-
- if (status == Ydb::StatusIds::OVERLOADED) {
- return ERetryErrorClass::LongRetry;
- }
-
- return ERetryErrorClass::NoRetry;
- }
-
- void CreateRetryState() {
- IRetryPolicy::TPtr policy = IRetryPolicy::GetExponentialBackoffPolicy(Retryable, TDuration::MilliSeconds(10), TDuration::MilliSeconds(200), TDuration::Seconds(1), std::numeric_limits<size_t>::max(), LeaseDuration / 2);
- RetryState = policy->CreateRetryState();
- }
-
private:
TActorId RunScriptActorId;
TString Database;
@@ -568,7 +511,6 @@ private:
TDuration LeaseDuration;
TIntrusivePtr<TKqpCounters> Counters;
TInstant LeaseUpdateStartTime;
- IRetryPolicy::IRetryState::TPtr RetryState = nullptr;
};
class TCheckLeaseStatusActorBase : public TActorBootstrapped<TCheckLeaseStatusActorBase> {
@@ -646,9 +588,9 @@ private:
}
WaitFinishQuery = true;
- FinalOperationStatus = ScriptFinalizeRequest->OperationStatus;
- FinalExecStatus = ScriptFinalizeRequest->ExecStatus;
- FinalIssues = ScriptFinalizeRequest->Issues;
+ FinalOperationStatus = ScriptFinalizeRequest->Description.OperationStatus;
+ FinalExecStatus = ScriptFinalizeRequest->Description.ExecStatus;
+ FinalIssues = ScriptFinalizeRequest->Description.Issues;
Send(MakeKqpFinalizeScriptServiceId(SelfId().NodeId()), ScriptFinalizeRequest.release());
}
@@ -1756,43 +1698,9 @@ private:
const TString SerializedMetas;
};
-class TSaveScriptExecutionResultMetaActor : public TActorBootstrapped<TSaveScriptExecutionResultMetaActor> {
-public:
- TSaveScriptExecutionResultMetaActor(const NActors::TActorId& replyActorId, const TString& database, const TString& executionId, const TString& serializedMetas)
- : ReplyActorId(replyActorId), Database(database), ExecutionId(executionId), SerializedMetas(serializedMetas)
- {
- }
-
- void Bootstrap() {
- Register(new TSaveScriptExecutionResultMetaQuery(Database, ExecutionId, SerializedMetas));
-
- Become(&TSaveScriptExecutionResultMetaActor::StateFunc);
- }
-
- STRICT_STFUNC(StateFunc,
- hFunc(TEvSaveScriptResultMetaFinished, Handle);
- )
-
- void Handle(TEvSaveScriptResultMetaFinished::TPtr& ev) {
- if (ev->Get()->Status == Ydb::StatusIds::ABORTED) {
- Register(new TSaveScriptExecutionResultMetaQuery(Database, ExecutionId, SerializedMetas));
- return;
- }
-
- Send(ev->Forward(ReplyActorId));
- PassAway();
- }
-
-private:
- const NActors::TActorId ReplyActorId;
- const TString Database;
- const TString ExecutionId;
- const TString SerializedMetas;
-};
-
class TSaveScriptExecutionResultQuery : public TQueryBase {
public:
- TSaveScriptExecutionResultQuery(const TString& database, const TString& executionId, i32 resultSetId, TMaybe<TInstant> expireAt, i64 firstRow, Ydb::ResultSet&& resultSet)
+ TSaveScriptExecutionResultQuery(const TString& database, const TString& executionId, i32 resultSetId, TMaybe<TInstant> expireAt, i64 firstRow, Ydb::ResultSet resultSet)
: Database(database), ExecutionId(executionId), ResultSetId(resultSetId), ExpireAt(expireAt), FirstRow(firstRow), ResultSet(std::move(resultSet))
{
}
@@ -1895,7 +1803,7 @@ public:
}
i64 numberRows = ResultSets.back().rows_size();
- Register(new TSaveScriptExecutionResultQuery(Database, ExecutionId, ResultSetId, ExpireAt, FirstRow, std::move(ResultSets.back())));
+ Register(new TQueryRetryActor<TSaveScriptExecutionResultQuery, TEvSaveScriptResultFinished, TString, TString, i32, TMaybe<TInstant>, i64, Ydb::ResultSet>(SelfId(), Database, ExecutionId, ResultSetId, ExpireAt, FirstRow, ResultSets.back()));
FirstRow += numberRows;
ResultSets.pop_back();
@@ -2203,8 +2111,8 @@ private:
class TSaveScriptExternalEffectActor : public TQueryBase {
public:
- explicit TSaveScriptExternalEffectActor(TEvSaveScriptExternalEffectRequest::TPtr ev)
- : Request(std::move(ev))
+ explicit TSaveScriptExternalEffectActor(const TEvSaveScriptExternalEffectRequest::TDescription& request)
+ : Request(request)
{}
void OnRunQuery() override {
@@ -2229,22 +2137,22 @@ public:
NYdb::TParamsBuilder params;
params
.AddParam("$database")
- .Utf8(Request->Get()->Database)
+ .Utf8(Request.Database)
.Build()
.AddParam("$execution_id")
- .Utf8(Request->Get()->ExecutionId)
+ .Utf8(Request.ExecutionId)
.Build()
.AddParam("$customer_supplied_id")
- .Utf8(Request->Get()->CustomerSuppliedId)
+ .Utf8(Request.CustomerSuppliedId)
.Build()
.AddParam("$user_token")
- .Utf8(Request->Get()->UserToken)
+ .Utf8(Request.UserToken)
.Build()
.AddParam("$script_sinks")
- .JsonDocument(SerializeSinks(Request->Get()->Sinks))
+ .JsonDocument(SerializeSinks(Request.Sinks))
.Build()
.AddParam("$script_secret_names")
- .JsonDocument(SerializeSecretNames(Request->Get()->SecretNames))
+ .JsonDocument(SerializeSecretNames(Request.SecretNames))
.Build();
RunDataQuery(sql, &params);
@@ -2255,7 +2163,7 @@ public:
}
void OnFinish(Ydb::StatusIds::StatusCode status, NYql::TIssues&& issues) override {
- Send(Request->Sender, new TEvSaveScriptExternalEffectResponse(status, std::move(issues)));
+ Send(Owner, new TEvSaveScriptExternalEffectResponse(status, std::move(issues)));
}
private:
@@ -2292,14 +2200,16 @@ private:
}
private:
- TEvSaveScriptExternalEffectRequest::TPtr Request;
+ TEvSaveScriptExternalEffectRequest::TDescription Request;
};
class TSaveScriptFinalStatusActor : public TQueryBase {
public:
- explicit TSaveScriptFinalStatusActor(TEvScriptFinalizeRequest::TPtr ev)
- : Request(ev)
- {}
+ explicit TSaveScriptFinalStatusActor(const TEvScriptFinalizeRequest::TDescription& request)
+ : Request(request)
+ {
+ Response = std::make_unique<TEvSaveScriptFinalStatusResponse>();
+ }
void OnRunQuery() override {
TString sql = R"(
@@ -2328,10 +2238,10 @@ public:
NYdb::TParamsBuilder params;
params
.AddParam("$database")
- .Utf8(Request->Get()->Database)
+ .Utf8(Request.Database)
.Build()
.AddParam("$execution_id")
- .Utf8(Request->Get()->ExecutionId)
+ .Utf8(Request.ExecutionId)
.Build();
RunDataQuery(sql, &params, TTxControl::BeginTx());
@@ -2354,16 +2264,16 @@ public:
TMaybe<i32> finalizationStatus = result.ColumnParser("finalization_status").GetOptionalInt32();
if (finalizationStatus) {
- if (Request->Get()->FinalizationStatus != *finalizationStatus) {
+ if (Request.FinalizationStatus != *finalizationStatus) {
Finish(Ydb::StatusIds::PRECONDITION_FAILED, "Execution already have different finalization status");
return;
}
- ApplicateScriptExternalEffectRequired = true;
+ Response->ApplicateScriptExternalEffectRequired = true;
}
TMaybe<i32> operationStatus = result.ColumnParser("operation_status").GetOptionalInt32();
- if (Request->Get()->LeaseGeneration && !operationStatus) {
+ if (Request.LeaseGeneration && !operationStatus) {
NYdb::TResultSetParser leaseResult(ResultSets[1]);
if (leaseResult.RowsCount() == 0) {
Finish(Ydb::StatusIds::INTERNAL_ERROR, "Unexpected operation state");
@@ -2378,7 +2288,7 @@ public:
return;
}
- if (*Request->Get()->LeaseGeneration != static_cast<ui64>(*leaseGenerationInDatabase)) {
+ if (*Request.LeaseGeneration != static_cast<ui64>(*leaseGenerationInDatabase)) {
Finish(Ydb::StatusIds::PRECONDITION_FAILED, "Lease was lost");
return;
}
@@ -2386,12 +2296,12 @@ public:
TMaybe<TString> customerSuppliedId = result.ColumnParser("customer_supplied_id").GetOptionalUtf8();
if (customerSuppliedId) {
- CustomerSuppliedId = *customerSuppliedId;
+ Response->CustomerSuppliedId = *customerSuppliedId;
}
TMaybe<TString> userToken = result.ColumnParser("user_token").GetOptionalUtf8();
if (userToken) {
- UserToken = *userToken;
+ Response->UserToken = *userToken;
}
SerializedSinks = result.ColumnParser("script_sinks").GetOptionalJsonDocument();
@@ -2408,7 +2318,7 @@ public:
NKqpProto::TKqpExternalSink sink;
NProtobufJson::Json2Proto(*serializedSink, sink);
- Sinks.push_back(sink);
+ Response->Sinks.push_back(sink);
}
}
@@ -2424,7 +2334,7 @@ public:
const NJson::TJsonValue* serializedSecretName;
value.GetValuePointer(i, &serializedSecretName);
- SecretNames.push_back(serializedSecretName->GetString());
+ Response->SecretNames.push_back(serializedSecretName->GetString());
}
}
@@ -2443,12 +2353,12 @@ public:
if (operationStatus) {
FinalStatusAlreadySaved = true;
- OperationAlreadyFinalized = !finalizationStatus;
+ Response->OperationAlreadyFinalized = !finalizationStatus;
CommitTransaction();
return;
}
- ApplicateScriptExternalEffectRequired = ApplicateScriptExternalEffectRequired || HasExternalEffect();
+ Response->ApplicateScriptExternalEffectRequired = Response->ApplicateScriptExternalEffectRequired || HasExternalEffect();
FinishScriptExecution();
}
@@ -2493,10 +2403,10 @@ public:
)";
TString serializedStats = "{}";
- if (Request->Get()->QueryStats) {
+ if (Request.QueryStats) {
NJson::TJsonValue statsJson;
Ydb::TableStats::QueryStats queryStats;
- NGRpcService::FillQueryStats(queryStats, *Request->Get()->QueryStats);
+ NGRpcService::FillQueryStats(queryStats, *Request.QueryStats);
NProtobufJson::Proto2Json(queryStats, statsJson, NProtobufJson::TProto2JsonConfig());
serializedStats = NJson::WriteJson(statsJson);
}
@@ -2504,40 +2414,40 @@ public:
NYdb::TParamsBuilder params;
params
.AddParam("$database")
- .Utf8(Request->Get()->Database)
+ .Utf8(Request.Database)
.Build()
.AddParam("$execution_id")
- .Utf8(Request->Get()->ExecutionId)
+ .Utf8(Request.ExecutionId)
.Build()
.AddParam("$operation_status")
- .Int32(Request->Get()->OperationStatus)
+ .Int32(Request.OperationStatus)
.Build()
.AddParam("$execution_status")
- .Int32(Request->Get()->ExecStatus)
+ .Int32(Request.ExecStatus)
.Build()
.AddParam("$finalization_status")
- .Int32(Request->Get()->FinalizationStatus)
+ .Int32(Request.FinalizationStatus)
.Build()
.AddParam("$issues")
- .JsonDocument(SerializeIssues(Request->Get()->Issues))
+ .JsonDocument(SerializeIssues(Request.Issues))
.Build()
.AddParam("$plan")
- .JsonDocument(Request->Get()->QueryPlan.value_or("{}"))
+ .JsonDocument(Request.QueryPlan.value_or("{}"))
.Build()
.AddParam("$stats")
.JsonDocument(serializedStats)
.Build()
.AddParam("$ast")
- .Utf8(Request->Get()->QueryAst.value_or(""))
+ .Utf8(Request.QueryAst.value_or(""))
.Build()
.AddParam("$operation_ttl")
.Interval(static_cast<i64>(OperationTtl.MicroSeconds()))
.Build()
.AddParam("$customer_supplied_id")
- .Utf8(CustomerSuppliedId)
+ .Utf8(Response->CustomerSuppliedId)
.Build()
.AddParam("$user_token")
- .Utf8(UserToken)
+ .Utf8(Response->UserToken)
.Build()
.AddParam("$script_sinks")
.OptionalJsonDocument(SerializedSinks)
@@ -2546,7 +2456,7 @@ public:
.OptionalJsonDocument(SerializedSecretNames)
.Build()
.AddParam("$applicate_script_external_effect_required")
- .Bool(ApplicateScriptExternalEffectRequired)
+ .Bool(Response->ApplicateScriptExternalEffectRequired)
.Build();
RunDataQuery(sql, &params, TTxControl::ContinueAndCommitTx());
@@ -2559,42 +2469,31 @@ public:
void OnFinish(Ydb::StatusIds::StatusCode status, NYql::TIssues&& issues) override {
if (!FinalStatusAlreadySaved) {
- KQP_PROXY_LOG_D("Finish script execution operation. ExecutionId: " << Request->Get()->ExecutionId
- << ". " << Ydb::StatusIds::StatusCode_Name(Request->Get()->OperationStatus)
- << ". Issues: " << Request->Get()->Issues.ToOneLineString() << ". Plan: " << Request->Get()->QueryPlan.value_or(""));
- }
-
- if (!ApplicateScriptExternalEffectRequired || status != Ydb::StatusIds::SUCCESS) {
- Send(Owner, new TEvScriptExecutionFinished(OperationAlreadyFinalized, status, issues));
- return;
+ KQP_PROXY_LOG_D("Finish script execution operation. ExecutionId: " << Request.ExecutionId
+ << ". " << Ydb::StatusIds::StatusCode_Name(Request.OperationStatus)
+ << ". Issues: " << Request.Issues.ToOneLineString() << ". Plan: " << Request.QueryPlan.value_or(""));
}
- auto response = std::make_unique<TEvSaveScriptFinalStatusResponse>(CustomerSuppliedId, UserToken);
- response->Sinks = std::move(Sinks);
- response->SecretNames = std::move(SecretNames);
+ Response->Status = status;
+ Response->Issues = std::move(issues);
- Send(Owner, response.release());
+ Send(Owner, Response.release());
}
private:
bool HasExternalEffect() const {
- return !Sinks.empty();
+ return !Response->Sinks.empty();
}
private:
- TEvScriptFinalizeRequest::TPtr Request;
+ TEvScriptFinalizeRequest::TDescription Request;
+ std::unique_ptr<TEvSaveScriptFinalStatusResponse> Response;
- bool OperationAlreadyFinalized = false;
bool FinalStatusAlreadySaved = false;
- bool ApplicateScriptExternalEffectRequired = false;
TDuration OperationTtl;
- TString CustomerSuppliedId;
- TString UserToken;
TMaybe<TString> SerializedSinks;
- std::vector<NKqpProto::TKqpExternalSink> Sinks;
TMaybe<TString> SerializedSecretNames;
- std::vector<TString> SecretNames;
};
class TScriptFinalizationFinisherActor : public TQueryBase {
@@ -2830,7 +2729,7 @@ NActors::IActor* CreateScriptLeaseUpdateActor(const TActorId& runScriptActorId,
}
NActors::IActor* CreateSaveScriptExecutionResultMetaActor(const NActors::TActorId& runScriptActorId, const TString& database, const TString& executionId, const TString& serializedMeta) {
- return new TSaveScriptExecutionResultMetaActor(runScriptActorId, database, executionId, serializedMeta);
+ return new TQueryRetryActor<TSaveScriptExecutionResultMetaQuery, TEvSaveScriptResultMetaFinished, TString, TString, TString>(runScriptActorId, database, executionId, serializedMeta);
}
NActors::IActor* CreateSaveScriptExecutionResultActor(const NActors::TActorId& runScriptActorId, const TString& database, const TString& executionId, i32 resultSetId, TMaybe<TInstant> expireAt, i64 firstRow, Ydb::ResultSet&& resultSet) {
@@ -2842,15 +2741,15 @@ NActors::IActor* CreateGetScriptExecutionResultActor(const NActors::TActorId& re
}
NActors::IActor* CreateSaveScriptExternalEffectActor(TEvSaveScriptExternalEffectRequest::TPtr ev) {
- return new TSaveScriptExternalEffectActor(std::move(ev));
+ return new TQueryRetryActor<TSaveScriptExternalEffectActor, TEvSaveScriptExternalEffectResponse, TEvSaveScriptExternalEffectRequest::TDescription>(ev->Sender, ev->Get()->Description);
}
-NActors::IActor* CreateSaveScriptFinalStatusActor(TEvScriptFinalizeRequest::TPtr ev) {
- return new TSaveScriptFinalStatusActor(std::move(ev));
+NActors::IActor* CreateSaveScriptFinalStatusActor(const NActors::TActorId& finalizationActorId, TEvScriptFinalizeRequest::TPtr ev) {
+ return new TQueryRetryActor<TSaveScriptFinalStatusActor, TEvSaveScriptFinalStatusResponse, TEvScriptFinalizeRequest::TDescription>(finalizationActorId, ev->Get()->Description);
}
-NActors::IActor* CreateScriptFinalizationFinisherActor(const TString& executionId, const TString& database, std::optional<Ydb::StatusIds::StatusCode> operationStatus, NYql::TIssues operationIssues) {
- return new TScriptFinalizationFinisherActor(executionId, database, operationStatus, std::move(operationIssues));
+NActors::IActor* CreateScriptFinalizationFinisherActor(const NActors::TActorId& finalizationActorId, const TString& executionId, const TString& database, std::optional<Ydb::StatusIds::StatusCode> operationStatus, NYql::TIssues operationIssues) {
+ return new TQueryRetryActor<TScriptFinalizationFinisherActor, TEvScriptExecutionFinished, TString, TString, std::optional<Ydb::StatusIds::StatusCode>, NYql::TIssues>(finalizationActorId, executionId, database, operationStatus, operationIssues);
}
NActors::IActor* CreateScriptProgressActor(const TString& executionId, const TString& database, const TString& queryPlan, const TString& queryStats) {
diff --git a/ydb/core/kqp/proxy_service/kqp_script_executions.h b/ydb/core/kqp/proxy_service/kqp_script_executions.h
index 5781046a1df..ea4fae00e84 100644
--- a/ydb/core/kqp/proxy_service/kqp_script_executions.h
+++ b/ydb/core/kqp/proxy_service/kqp_script_executions.h
@@ -33,8 +33,8 @@ NActors::IActor* CreateGetScriptExecutionResultActor(const NActors::TActorId& re
// Compute external effects and updates status in database
NActors::IActor* CreateSaveScriptExternalEffectActor(TEvSaveScriptExternalEffectRequest::TPtr ev);
-NActors::IActor* CreateSaveScriptFinalStatusActor(TEvScriptFinalizeRequest::TPtr ev);
-NActors::IActor* CreateScriptFinalizationFinisherActor(const TString& executionId, const TString& database, std::optional<Ydb::StatusIds::StatusCode> operationStatus, NYql::TIssues operationIssues);
+NActors::IActor* CreateSaveScriptFinalStatusActor(const NActors::TActorId& finalizationActorId, TEvScriptFinalizeRequest::TPtr ev);
+NActors::IActor* CreateScriptFinalizationFinisherActor(const NActors::TActorId& finalizationActorId, const TString& executionId, const TString& database, std::optional<Ydb::StatusIds::StatusCode> operationStatus, NYql::TIssues operationIssues);
NActors::IActor* CreateScriptProgressActor(const TString& executionId, const TString& database, const TString& queryPlan, const TString& queryStats);
} // namespace NKikimr::NKqp
diff --git a/ydb/library/query_actor/query_actor.h b/ydb/library/query_actor/query_actor.h
index ef47d2300a0..5d21f2f840e 100644
--- a/ydb/library/query_actor/query_actor.h
+++ b/ydb/library/query_actor/query_actor.h
@@ -12,11 +12,12 @@
#include <ydb/library/actors/core/actorsystem.h>
#include <ydb/library/actors/core/event_local.h>
#include <ydb/library/actors/core/events.h>
+#include <ydb/library/actors/core/hfunc.h>
+#include <library/cpp/retry/retry_policy.h>
#include <library/cpp/threading/future/future.h>
namespace NKikimr {
-// TODO: add retry logic
class TQueryBase : public NActors::TActorBootstrapped<TQueryBase> {
protected:
struct TTxControl {
@@ -168,4 +169,101 @@ protected:
std::vector<NYdb::TResultSet> ResultSets;
};
+template<typename TQueryActor, typename TResponse, typename ...TArgs>
+class TQueryRetryActor : public NActors::TActorBootstrapped<TQueryRetryActor<TQueryActor, TResponse, TArgs...>> {
+public:
+ using TBase = NActors::TActorBootstrapped<TQueryRetryActor<TQueryActor, TResponse, TArgs...>>;
+ using IRetryPolicy = IRetryPolicy<Ydb::StatusIds::StatusCode>;
+
+ explicit TQueryRetryActor(const NActors::TActorId& replyActorId, const TArgs&... args, TDuration maxRetryTime = TDuration::Seconds(1))
+ : ReplyActorId(replyActorId)
+ , RetryPolicy(IRetryPolicy::GetExponentialBackoffPolicy(
+ Retryable, TDuration::MilliSeconds(10),
+ TDuration::MilliSeconds(200), TDuration::Seconds(1),
+ std::numeric_limits<size_t>::max(), maxRetryTime
+ ))
+ , CreateQueryActor([=]() {
+ return new TQueryActor(args...);
+ })
+ {}
+
+ TQueryRetryActor(const NActors::TActorId& replyActorId, IRetryPolicy::TPtr retryPolicy, const TArgs&... args)
+ : ReplyActorId(replyActorId)
+ , RetryPolicy(retryPolicy)
+ , CreateQueryActor([=]() {
+ return new TQueryActor(args...);
+ })
+ {}
+
+ void StartQueryActor() const {
+ TBase::Register(CreateQueryActor());
+ }
+
+ void Bootstrap() {
+ TBase::Become(&TQueryRetryActor::StateFunc);
+ StartQueryActor();
+ }
+
+ STRICT_STFUNC(StateFunc,
+ hFunc(NActors::TEvents::TEvWakeup, Wakeup);
+ hFunc(TResponse, Handle);
+ )
+
+ void Wakeup(NActors::TEvents::TEvWakeup::TPtr&) {
+ StartQueryActor();
+ }
+
+ void Handle(const typename TResponse::TPtr& ev) {
+ const Ydb::StatusIds::StatusCode status = ev->Get()->Status;
+ if (Retryable(status) == ERetryErrorClass::NoRetry) {
+ Reply(ev);
+ return;
+ }
+
+ if (RetryState == nullptr) {
+ RetryState = RetryPolicy->CreateRetryState();
+ }
+
+ if (auto delay = RetryState->GetNextRetryDelay(status)) {
+ TBase::Schedule(*delay, new NActors::TEvents::TEvWakeup());
+ } else {
+ Reply(ev);
+ }
+ }
+
+ void Reply(const typename TResponse::TPtr& ev) {
+ TBase::Send(ev->Forward(ReplyActorId));
+ TBase::PassAway();
+ }
+
+ static ERetryErrorClass Retryable(Ydb::StatusIds::StatusCode status) {
+ if (status == Ydb::StatusIds::SUCCESS) {
+ return ERetryErrorClass::NoRetry;
+ }
+
+ if (status == Ydb::StatusIds::INTERNAL_ERROR
+ || status == Ydb::StatusIds::UNAVAILABLE
+ || status == Ydb::StatusIds::BAD_SESSION
+ || status == Ydb::StatusIds::SESSION_EXPIRED
+ || status == Ydb::StatusIds::SESSION_BUSY
+ || status == Ydb::StatusIds::TIMEOUT
+ || status == Ydb::StatusIds::ABORTED) {
+ return ERetryErrorClass::ShortRetry;
+ }
+
+ if (status == Ydb::StatusIds::OVERLOADED
+ || status == Ydb::StatusIds::UNDETERMINED) {
+ return ERetryErrorClass::LongRetry;
+ }
+
+ return ERetryErrorClass::NoRetry;
+ }
+
+private:
+ const NActors::TActorId ReplyActorId;
+ const IRetryPolicy::TPtr RetryPolicy;
+ const std::function<TQueryActor*()> CreateQueryActor;
+ IRetryPolicy::IRetryState::TPtr RetryState = nullptr;
+};
+
} // namespace NKikimr