diff options
| author | Pisarenko Grigoriy <[email protected]> | 2024-02-07 00:21:44 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-02-07 00:21:44 +0300 |
| commit | d010ed3f89cde39cb409fbe2448e871bced2eb43 (patch) | |
| tree | 17934e9d2122a9defaa4ef0c842e648a918cd847 | |
| parent | e889ee1e6fc13f3ea314b24cf412f371fb09b9da (diff) | |
YQ-2734 added retries for internal queries (#1457)
* Added retries for internal queries
* Moved logic to query_actor.h
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, ¶ms); @@ -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, ¶ms, 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, ¶ms, 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 |
