diff options
| author | Hor911 <[email protected]> | 2024-03-12 15:18:54 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-03-12 15:18:54 +0300 |
| commit | df862affe0123b3bd96f0950822d8d1a57ddf2f7 (patch) | |
| tree | 9a9d64b3b0432c8bf4a47449c4bfd9bbea8f1291 | |
| parent | bfc374db311ff3ccd9a0bfeb65bca7ba15706045 (diff) | |
Hint-based statistics (#2647)
| -rw-r--r-- | ydb/core/fq/libs/compute/common/utils.cpp | 144 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/common/utils.h | 33 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/actors_factory.cpp | 44 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/base_status_updater_actor.h | 96 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/events/events.h | 4 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/executer_actor.cpp | 9 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/executer_actor.h | 1 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp | 131 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/status_tracker_actor.h | 2 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/stopper_actor.cpp | 42 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/stopper_actor.h | 2 | ||||
| -rw-r--r-- | ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp | 10 | ||||
| -rw-r--r-- | ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp | 3 | ||||
| -rw-r--r-- | ydb/core/fq/libs/protos/fq_private.proto | 1 |
14 files changed, 337 insertions, 185 deletions
diff --git a/ydb/core/fq/libs/compute/common/utils.cpp b/ydb/core/fq/libs/compute/common/utils.cpp index 127bd26f1ec..b790d8d8a8b 100644 --- a/ydb/core/fq/libs/compute/common/utils.cpp +++ b/ydb/core/fq/libs/compute/common/utils.cpp @@ -1,8 +1,11 @@ #include "utils.h" #include <library/cpp/json/json_reader.h> +#include <library/cpp/json/json_writer.h> #include <library/cpp/json/yson/json2yson.h> +#include <ydb/core/fq/libs/control_plane_storage/internal/utils.h> + namespace NFq { using TAggregates = std::map<TString, std::optional<ui64>>; @@ -621,7 +624,7 @@ void EnumeratePlansV2(NYson::TYsonWriter& writer, NJson::TJsonValue& value, ui32 } } -TString GetV1StatFromV2PlanV2(const TString& plan) { +TString GetV1StatFromV2PlanV2(const TString& plan, double* cpuUsage) { TStringStream out; NYson::TYsonWriter writer(&out); writer.OnBeginMap(); @@ -655,6 +658,9 @@ TString GetV1StatFromV2PlanV2(const TString& plan) { if (totals.CpuTimeUs.Sum) { writer.OnKeyedItem("cpu"); writer.OnStringScalar(FormatDurationUs(totals.CpuTimeUs.Sum)); + if (cpuUsage) { + *cpuUsage = totals.CpuTimeUs.Sum / 1000000.0; + } } if (totals.SourceCpuTimeUs.Sum) { writer.OnKeyedItem("scpu"); @@ -750,4 +756,140 @@ TPublicStat GetPublicStat(const TString& statistics) { return counters; } +struct TNoneStatProcessor : IPlanStatProcessor { + Ydb::Query::StatsMode GetStatsMode() override { + return Ydb::Query::StatsMode::STATS_MODE_NONE; + } + + TString ConvertPlan(TString& plan) override { + return plan; + } + + TString GetQueryStat(TString&, double& cpuUsage) override { + cpuUsage = 0.0; + return ""; + } + + TPublicStat GetPublicStat(TString&) override { + return TPublicStat{}; + } +}; + +struct TBasicStatProcessor : TNoneStatProcessor { + Ydb::Query::StatsMode GetStatsMode() override { + return Ydb::Query::StatsMode::STATS_MODE_BASIC; + } +}; + +struct TFullStatProcessor : IPlanStatProcessor { + Ydb::Query::StatsMode GetStatsMode() override { + return Ydb::Query::StatsMode::STATS_MODE_FULL; + } + + TString ConvertPlan(TString& plan) override { + return plan; + } + + TString GetQueryStat(TString& plan, double& cpuUsage) override { + return GetV1StatFromV2Plan(plan, &cpuUsage); + } + + TPublicStat GetPublicStat(TString& stat) override { + return NFq::GetPublicStat(stat); + } +}; + +struct TProfileStatProcessor : TFullStatProcessor { + Ydb::Query::StatsMode GetStatsMode() override { + return Ydb::Query::StatsMode::STATS_MODE_PROFILE; + } +}; + +struct TProdStatProcessor : TFullStatProcessor { + TString GetQueryStat(TString& plan, double& cpuUsage) override { + return GetPrettyStatistics(GetV1StatFromV2Plan(plan, &cpuUsage)); + } +}; + +std::unique_ptr<IPlanStatProcessor> CreateStatProcessor(const TString& statViewName) { + // disallow none and basic stat since they do not support metering + // if (statViewName == "stat_none") return std::make_unique<TNoneStatProcessor>(); + // if (statViewName == "stat_basc") return std::make_unique<TBasicStatProcessor>(); + if (statViewName == "stat_full") return std::make_unique<TFullStatProcessor>(); + if (statViewName == "stat_prof") return std::make_unique<TProfileStatProcessor>(); + if (statViewName == "stat_prod") return std::make_unique<TProdStatProcessor>(); + return std::make_unique<TFullStatProcessor>(); +} + +PingTaskRequestBuilder::PingTaskRequestBuilder(const NConfig::TCommonConfig& commonConfig, std::unique_ptr<IPlanStatProcessor>&& processor) + : Compressor(commonConfig.GetQueryArtifactsCompressionMethod(), commonConfig.GetQueryArtifactsCompressionMinSize()) + , Processor(std::move(processor)) +{} + +Fq::Private::PingTaskRequest PingTaskRequestBuilder::Build( + const Ydb::TableStats::QueryStats& queryStats, + const NYql::TIssues& issues, + std::optional<FederatedQuery::QueryMeta::ComputeStatus> computeStatus, + std::optional<NYql::NDqProto::StatusIds::StatusCode> pendingStatusCode +) { + Fq::Private::PingTaskRequest pingTaskRequest = Build(queryStats); + + if (issues) { + NYql::IssuesToMessage(issues, pingTaskRequest.mutable_issues()); + } + + if (computeStatus) { + pingTaskRequest.set_status(*computeStatus); + } + + if (pendingStatusCode) { + pingTaskRequest.set_pending_status_code(*pendingStatusCode); + } + + return pingTaskRequest; +} + + +Fq::Private::PingTaskRequest PingTaskRequestBuilder::Build(const Ydb::TableStats::QueryStats& queryStats) { + return Build(queryStats.query_plan(), queryStats.query_ast()); +} + +Fq::Private::PingTaskRequest PingTaskRequestBuilder::Build(const TString& queryPlan, const TString& queryAst) { + Fq::Private::PingTaskRequest pingTaskRequest; + + Issues.Clear(); + + auto plan = queryPlan; + try { + plan = Processor->ConvertPlan(plan); + } catch(const NJson::TJsonException& ex) { + Issues.AddIssue(NYql::TIssue(TStringBuilder() << "Error plan conversion: " << ex.what())); + } + + if (Compressor.IsEnabled()) { + auto [astCompressionMethod, astCompressed] = Compressor.Compress(queryAst); + pingTaskRequest.mutable_ast_compressed()->set_method(astCompressionMethod); + pingTaskRequest.mutable_ast_compressed()->set_data(astCompressed); + + auto [planCompressionMethod, planCompressed] = Compressor.Compress(plan); + pingTaskRequest.mutable_plan_compressed()->set_method(planCompressionMethod); + pingTaskRequest.mutable_plan_compressed()->set_data(planCompressed); + } else { + pingTaskRequest.set_ast(queryAst); + pingTaskRequest.set_plan(plan); + } + + CpuUsage = 0.0; + try { + auto stat = Processor->GetQueryStat(plan, CpuUsage); + pingTaskRequest.set_statistics(stat); + pingTaskRequest.set_dump_raw_statistics(true); + PublicStat = Processor->GetPublicStat(stat); + } catch(const NJson::TJsonException& ex) { + Issues.AddIssue(NYql::TIssue(TStringBuilder() << "Error stat conversion: " << ex.what())); + } + + return pingTaskRequest; +} + } // namespace NFq diff --git a/ydb/core/fq/libs/compute/common/utils.h b/ydb/core/fq/libs/compute/common/utils.h index 4a61a45bf61..47387490162 100644 --- a/ydb/core/fq/libs/compute/common/utils.h +++ b/ydb/core/fq/libs/compute/common/utils.h @@ -1,8 +1,12 @@ #pragma once +#include <memory> + +#include <ydb/core/fq/libs/common/compression.h> #include <ydb/core/fq/libs/compute/common/config.h> #include <ydb/core/fq/libs/shared_resources/shared_resources.h> #include <ydb/core/fq/libs/ydb/ydb.h> + #include <ydb/public/sdk/cpp/client/ydb_table/table.h> namespace NFq { @@ -43,4 +47,33 @@ struct TPublicStat { TPublicStat GetPublicStat(const TString& statistics); +struct IPlanStatProcessor { + virtual ~IPlanStatProcessor() = default; + virtual Ydb::Query::StatsMode GetStatsMode() = 0; + virtual TString ConvertPlan(TString& plan) = 0; + virtual TString GetQueryStat(TString& plan, double& cpuUsage) = 0; + virtual TPublicStat GetPublicStat(TString& stat) = 0; +}; + +std::unique_ptr<IPlanStatProcessor> CreateStatProcessor(const TString& statViewName); + +class PingTaskRequestBuilder { +public: + PingTaskRequestBuilder(const NConfig::TCommonConfig& commonConfig, std::unique_ptr<IPlanStatProcessor>&& processor); + Fq::Private::PingTaskRequest Build( + const Ydb::TableStats::QueryStats& queryStats, + const NYql::TIssues& issues, + std::optional<FederatedQuery::QueryMeta::ComputeStatus> computeStatus = std::nullopt, + std::optional<NYql::NDqProto::StatusIds::StatusCode> pendingStatusCode = std::nullopt + ); + Fq::Private::PingTaskRequest Build(const Ydb::TableStats::QueryStats& queryStats); + Fq::Private::PingTaskRequest Build(const TString& queryPlan, const TString& queryAst); + NYql::TIssues Issues; + double CpuUsage = 0.0; + TPublicStat PublicStat; +private: + const TCompressor Compressor; + std::unique_ptr<IPlanStatProcessor> Processor; +}; + } // namespace NFq diff --git a/ydb/core/fq/libs/compute/ydb/actors_factory.cpp b/ydb/core/fq/libs/compute/ydb/actors_factory.cpp index 6dd8ea66a3a..aa7d38d00fc 100644 --- a/ydb/core/fq/libs/compute/ydb/actors_factory.cpp +++ b/ydb/core/fq/libs/compute/ydb/actors_factory.cpp @@ -9,6 +9,7 @@ #include "ydb_connector_actor.h" #include <ydb/core/fq/libs/compute/common/pinger.h> +#include <ydb/core/fq/libs/compute/common/utils.h> namespace NFq { @@ -16,6 +17,7 @@ struct TActorFactory : public IActorFactory { TActorFactory(const NFq::TRunActorParams& params, const ::NYql::NCommon::TServiceCounters& counters) : Params(params) , Counters(counters) + , StatViewName(GetStatViewName()) {} std::unique_ptr<NActors::IActor> CreatePinger(const NActors::TActorId& parent) const override { @@ -46,14 +48,14 @@ struct TActorFactory : public IActorFactory { std::unique_ptr<NActors::IActor> CreateExecuter(const NActors::TActorId &parent, const NActors::TActorId &connector, const NActors::TActorId &pinger) const override { - return CreateExecuterActor(Params, parent, connector, pinger, Counters); + return CreateExecuterActor(Params, CreateStatProcessor()->GetStatsMode(), parent, connector, pinger, Counters); } std::unique_ptr<NActors::IActor> CreateStatusTracker(const NActors::TActorId &parent, const NActors::TActorId &connector, const NActors::TActorId &pinger, const NYdb::TOperation::TOperationId& operationId) const override { - return CreateStatusTrackerActor(Params, parent, connector, pinger, operationId, Counters); + return CreateStatusTrackerActor(Params, parent, connector, pinger, operationId, CreateStatProcessor(), Counters); } std::unique_ptr<NActors::IActor> CreateResultWriter(const NActors::TActorId& parent, @@ -82,12 +84,48 @@ struct TActorFactory : public IActorFactory { const NActors::TActorId& connector, const NActors::TActorId& pinger, const NYdb::TOperation::TOperationId& operationId) const override { - return CreateStopperActor(Params, parent, connector, pinger, operationId, Counters); + return CreateStopperActor(Params, parent, connector, pinger, operationId, CreateStatProcessor(), Counters); + } + + std::unique_ptr<IPlanStatProcessor> CreateStatProcessor() const { + return NFq::CreateStatProcessor(StatViewName); + } + + TString GetStatViewName() { + auto p = Params.Sql.find("--fq_dev_hint_"); + if (p != Params.Sql.npos) { + p += 14; + auto p1 = Params.Sql.find("\n", p); + TString mode = Params.Sql.substr(p, p1 == Params.Sql.npos ? Params.Sql.npos : p1 - p); + if (mode) { + return mode; + } + } + + if (!Params.Config.GetControlPlaneStorage().GetDumpRawStatistics()) { + return "stat_prod"; + } + + switch (Params.Config.GetControlPlaneStorage().GetStatsMode()) { + case Ydb::Query::StatsMode::STATS_MODE_UNSPECIFIED: + return "stat_full"; + case Ydb::Query::StatsMode::STATS_MODE_NONE: + return "stat_none"; + case Ydb::Query::StatsMode::STATS_MODE_BASIC: + return "stat_basc"; + case Ydb::Query::StatsMode::STATS_MODE_FULL: + return "stat_full"; + case Ydb::Query::StatsMode::STATS_MODE_PROFILE: + return "stat_prof"; + default: + return "stat_full"; + } } private: NFq::TRunActorParams Params; ::NYql::NCommon::TServiceCounters Counters; + TString StatViewName; }; IActorFactory::TPtr CreateActorFactory(const NFq::TRunActorParams& params, const ::NYql::NCommon::TServiceCounters& counters) { diff --git a/ydb/core/fq/libs/compute/ydb/base_status_updater_actor.h b/ydb/core/fq/libs/compute/ydb/base_status_updater_actor.h deleted file mode 100644 index d11de64d43c..00000000000 --- a/ydb/core/fq/libs/compute/ydb/base_status_updater_actor.h +++ /dev/null @@ -1,96 +0,0 @@ -#pragma once - -#include "base_compute_actor.h" - -#include <ydb/core/fq/libs/common/compression.h> -#include <ydb/core/fq/libs/compute/common/utils.h> - -#include <ydb/library/yql/public/issue/yql_issue_message.h> - -namespace NFq { - -template<typename TDerived> -class TBaseStatusUpdaterActor : public TBaseComputeActor<TDerived> { -public: - using TBase = TBaseComputeActor<TDerived>; - - TBaseStatusUpdaterActor(const NConfig::TCommonConfig& commonConfig, const ::NYql::NCommon::TServiceCounters& queryCounters, const TString& stepName) - : TBase(queryCounters, stepName) - , Compressor(commonConfig.GetQueryArtifactsCompressionMethod(), commonConfig.GetQueryArtifactsCompressionMinSize()) - {} - - TBaseStatusUpdaterActor(const NConfig::TCommonConfig& commonConfig, const ::NMonitoring::TDynamicCounterPtr& baseCounters, const TString& stepName) - : TBase(baseCounters, stepName) - , Compressor(commonConfig.GetQueryArtifactsCompressionMethod(), commonConfig.GetQueryArtifactsCompressionMinSize()) - {} - - void SetPingCounters(TComputeRequestCountersPtr pingCounters) { - PingCounters = std::move(pingCounters); - } - - void OnPingRequestStart() { - if (!PingCounters) { - return; - } - - StartTime = TInstant::Now(); - PingCounters->InFly->Inc(); - } - - void OnPingRequestFinish(bool success) { - if (!PingCounters) { - return; - } - - PingCounters->InFly->Dec(); - PingCounters->LatencyMs->Collect((TInstant::Now() - StartTime).MilliSeconds()); - if (success) { - PingCounters->Ok->Inc(); - } else { - PingCounters->Error->Inc(); - } - } - - Fq::Private::PingTaskRequest GetPingTaskRequest(std::optional<FederatedQuery::QueryMeta::ComputeStatus> computeStatus, std::optional<NYql::NDqProto::StatusIds::StatusCode> pendingStatusCode, const NYql::TIssues& issues, const Ydb::TableStats::QueryStats& queryStats) const { - Fq::Private::PingTaskRequest pingTaskRequest; - NYql::IssuesToMessage(issues, pingTaskRequest.mutable_issues()); - if (computeStatus) { - pingTaskRequest.set_status(*computeStatus); - } - if (pendingStatusCode) { - pingTaskRequest.set_pending_status_code(*pendingStatusCode); - } - PrepareAstAndPlan(pingTaskRequest, queryStats.query_plan(), queryStats.query_ast()); - return pingTaskRequest; - } - - // Can throw errors - Fq::Private::PingTaskRequest GetPingTaskRequestStatistics(std::optional<FederatedQuery::QueryMeta::ComputeStatus> computeStatus, std::optional<NYql::NDqProto::StatusIds::StatusCode> pendingStatusCode, const NYql::TIssues& issues, const Ydb::TableStats::QueryStats& queryStats, double* cpuUsage = nullptr) const { - Fq::Private::PingTaskRequest pingTaskRequest = GetPingTaskRequest(computeStatus, pendingStatusCode, issues, queryStats); - pingTaskRequest.set_statistics(GetV1StatFromV2Plan(queryStats.query_plan(), cpuUsage)); - return pingTaskRequest; - } - -protected: - void PrepareAstAndPlan(Fq::Private::PingTaskRequest& request, const TString& plan, const TString& expr) const { - if (Compressor.IsEnabled()) { - auto [astCompressionMethod, astCompressed] = Compressor.Compress(expr); - request.mutable_ast_compressed()->set_method(astCompressionMethod); - request.mutable_ast_compressed()->set_data(astCompressed); - - auto [planCompressionMethod, planCompressed] = Compressor.Compress(plan); - request.mutable_plan_compressed()->set_method(planCompressionMethod); - request.mutable_plan_compressed()->set_data(planCompressed); - } else { - request.set_ast(expr); - request.set_plan(plan); - } - } - -private: - TInstant StartTime; - TComputeRequestCountersPtr PingCounters; - const TCompressor Compressor; -}; - -} /* NFq */ diff --git a/ydb/core/fq/libs/compute/ydb/events/events.h b/ydb/core/fq/libs/compute/ydb/events/events.h index 442ff42e2cf..b893fa3d2d5 100644 --- a/ydb/core/fq/libs/compute/ydb/events/events.h +++ b/ydb/core/fq/libs/compute/ydb/events/events.h @@ -71,13 +71,14 @@ struct TEvYdbCompute { // Events struct TEvExecuteScriptRequest : public NActors::TEventLocal<TEvExecuteScriptRequest, EvExecuteScriptRequest> { - TEvExecuteScriptRequest(TString sql, TString idempotencyKey, const TDuration& resultTtl, const TDuration& operationTimeout, Ydb::Query::Syntax syntax, Ydb::Query::ExecMode execMode, const TString& traceId) + TEvExecuteScriptRequest(TString sql, TString idempotencyKey, const TDuration& resultTtl, const TDuration& operationTimeout, Ydb::Query::Syntax syntax, Ydb::Query::ExecMode execMode, Ydb::Query::StatsMode statsMode, const TString& traceId) : Sql(std::move(sql)) , IdempotencyKey(std::move(idempotencyKey)) , ResultTtl(resultTtl) , OperationTimeout(operationTimeout) , Syntax(syntax) , ExecMode(execMode) + , StatsMode(statsMode) , TraceId(traceId) {} @@ -87,6 +88,7 @@ struct TEvYdbCompute { TDuration OperationTimeout; Ydb::Query::Syntax Syntax = Ydb::Query::SYNTAX_YQL_V1; Ydb::Query::ExecMode ExecMode = Ydb::Query::EXEC_MODE_EXECUTE; + Ydb::Query::StatsMode StatsMode = Ydb::Query::StatsMode::STATS_MODE_FULL; TString TraceId; }; diff --git a/ydb/core/fq/libs/compute/ydb/executer_actor.cpp b/ydb/core/fq/libs/compute/ydb/executer_actor.cpp index 73a90ba51c6..177fe00ded3 100644 --- a/ydb/core/fq/libs/compute/ydb/executer_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/executer_actor.cpp @@ -59,9 +59,10 @@ public: } }; - TExecuterActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const ::NYql::NCommon::TServiceCounters& queryCounters) + TExecuterActor(const TRunActorParams& params, Ydb::Query::StatsMode statsMode, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const ::NYql::NCommon::TServiceCounters& queryCounters) : TBaseComputeActor(queryCounters, "Executer") , Params(params) + , StatsMode(statsMode) , Parent(parent) , Connector(connector) , Pinger(pinger) @@ -114,7 +115,7 @@ public: } void SendExecuteScript() { - Register(new TRetryActor<TEvYdbCompute::TEvExecuteScriptRequest, TEvYdbCompute::TEvExecuteScriptResponse, TString, TString, TDuration, TDuration, Ydb::Query::Syntax, Ydb::Query::ExecMode, TString>(Counters.GetCounters(ERequestType::RT_EXECUTE_SCRIPT), SelfId(), Connector, Params.Sql, Params.JobId, Params.ResultTtl, Params.ExecutionTtl, GetSyntax(), GetExecuteMode(), Params.JobId + "_" + ToString(Params.RestartCount))); + Register(new TRetryActor<TEvYdbCompute::TEvExecuteScriptRequest, TEvYdbCompute::TEvExecuteScriptResponse, TString, TString, TDuration, TDuration, Ydb::Query::Syntax, Ydb::Query::ExecMode, Ydb::Query::StatsMode, TString>(Counters.GetCounters(ERequestType::RT_EXECUTE_SCRIPT), SelfId(), Connector, Params.Sql, Params.JobId, Params.ResultTtl, Params.ExecutionTtl, GetSyntax(), GetExecuteMode(), StatsMode, Params.JobId + "_" + ToString(Params.RestartCount))); } Ydb::Query::Syntax GetSyntax() const { @@ -162,6 +163,7 @@ public: private: TRunActorParams Params; + Ydb::Query::StatsMode StatsMode; TActorId Parent; TActorId Connector; TActorId Pinger; @@ -172,11 +174,12 @@ private: }; std::unique_ptr<NActors::IActor> CreateExecuterActor(const TRunActorParams& params, + Ydb::Query::StatsMode statsMode, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const ::NYql::NCommon::TServiceCounters& queryCounters) { - return std::make_unique<TExecuterActor>(params, parent, connector, pinger, queryCounters); + return std::make_unique<TExecuterActor>(params, statsMode, parent, connector, pinger, queryCounters); } } diff --git a/ydb/core/fq/libs/compute/ydb/executer_actor.h b/ydb/core/fq/libs/compute/ydb/executer_actor.h index 76350109248..c1a6c1d6478 100644 --- a/ydb/core/fq/libs/compute/ydb/executer_actor.h +++ b/ydb/core/fq/libs/compute/ydb/executer_actor.h @@ -9,6 +9,7 @@ namespace NFq { std::unique_ptr<NActors::IActor> CreateExecuterActor(const TRunActorParams& params, + Ydb::Query::StatsMode statsMode, const NActors::TActorId& parent, const NActors::TActorId& connector, const NActors::TActorId& pinger, diff --git a/ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp b/ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp index 6b1d4fe0d3a..6b5c7584d03 100644 --- a/ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp @@ -1,4 +1,5 @@ -#include "base_status_updater_actor.h" +#include "base_compute_actor.h" +#include "status_tracker_actor.h" #include <ydb/core/fq/libs/common/util.h> #include <ydb/core/fq/libs/compute/common/metrics.h> @@ -34,10 +35,12 @@ namespace NFq { using namespace NActors; using namespace NFq; -class TStatusTrackerActor : public TBaseStatusUpdaterActor<TStatusTrackerActor> { +class TStatusTrackerActor : public TBaseComputeActor<TStatusTrackerActor> { public: using IRetryPolicy = IRetryPolicy<const TEvYdbCompute::TEvGetOperationResponse::TPtr&>; + using TBase = TBaseComputeActor<TStatusTrackerActor>; + enum ERequestType { RT_GET_OPERATION, RT_PING, @@ -66,18 +69,17 @@ public: } }; - TStatusTrackerActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, const ::NYql::NCommon::TServiceCounters& queryCounters) - : TBaseStatusUpdaterActor(params.Config.GetCommon(), queryCounters, "StatusTracker") + TStatusTrackerActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters) + : TBase(queryCounters, "StatusTracker") , Params(params) , Parent(parent) , Connector(connector) , Pinger(pinger) , OperationId(operationId) + , Builder(params.Config.GetCommon(), std::move(processor)) , Counters(GetStepCountersSubgroup()) , BackoffTimer(20, 1000) - { - SetPingCounters(Counters.GetCounters(ERequestType::RT_PING)); - } + {} static constexpr char ActorName[] = "FQ_STATUS_TRACKER"; @@ -93,7 +95,15 @@ public: ) void Handle(const TEvents::TEvForwardPingResponse::TPtr& ev) { - OnPingRequestFinish(ev.Get()->Get()->Success); + auto pingCounters = Counters.GetCounters(ERequestType::RT_PING); + pingCounters->InFly->Dec(); + pingCounters->LatencyMs->Collect((TInstant::Now() - StartTime).MilliSeconds()); + + if (ev.Get()->Get()->Success) { + pingCounters->Ok->Inc(); + } else { + pingCounters->Error->Inc(); + } if (ev->Cookie) { return; @@ -133,7 +143,6 @@ public: return; } - ReportPublicCounters(response.QueryStats); LOG_D("Execution status: " << static_cast<int>(response.ExecStatus)); switch (response.ExecStatus) { case NYdb::NQuery::EExecStatus::Unspecified: @@ -162,47 +171,42 @@ public: } } - void ReportPublicCounters(const Ydb::TableStats::QueryStats& stats) { - try { - auto stat = GetPublicStat(GetV1StatFromV2Plan(stats.query_plan())); - auto publicCounters = GetPublicCounters(); + void ReportPublicCounters(const TPublicStat& stat) { + auto publicCounters = GetPublicCounters(); - if (stat.MemoryUsageBytes) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.memory_usage_bytes"); - counter = *stat.MemoryUsageBytes; - } + if (stat.MemoryUsageBytes) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.memory_usage_bytes"); + counter = *stat.MemoryUsageBytes; + } - if (stat.CpuUsageUs) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.cpu_usage_us", true); - counter = *stat.CpuUsageUs; - } + if (stat.CpuUsageUs) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.cpu_usage_us", true); + counter = *stat.CpuUsageUs; + } - if (stat.InputBytes) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.input_bytes", true); - counter = *stat.InputBytes; - } + if (stat.InputBytes) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.input_bytes", true); + counter = *stat.InputBytes; + } - if (stat.OutputBytes) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.output_bytes", true); - counter = *stat.OutputBytes; - } + if (stat.OutputBytes) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.output_bytes", true); + counter = *stat.OutputBytes; + } - if (stat.SourceInputRecords) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.source_input_records", true); - counter = *stat.SourceInputRecords; - } + if (stat.SourceInputRecords) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.source_input_records", true); + counter = *stat.SourceInputRecords; + } - if (stat.SinkOutputRecords) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.sink_output_records", true); - counter = *stat.SinkOutputRecords; - } + if (stat.SinkOutputRecords) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.sink_output_records", true); + counter = *stat.SinkOutputRecords; + } - if (stat.RunningTasks) { - auto& counter = *publicCounters->GetNamedCounter("name", "query.running_tasks"); - counter = *stat.RunningTasks; - } - } catch(const NJson::TJsonException& ex) { - LOG_E("Error statistics conversion: " << ex.what()); + if (stat.RunningTasks) { + auto& counter = *publicCounters->GetNamedCounter("name", "query.running_tasks"); + counter = *stat.RunningTasks; } } @@ -210,22 +214,20 @@ public: Register(new TRetryActor<TEvYdbCompute::TEvGetOperationRequest, TEvYdbCompute::TEvGetOperationResponse, NYdb::TOperation::TOperationId>(Counters.GetCounters(ERequestType::RT_GET_OPERATION), delay, SelfId(), Connector, OperationId)); } - std::pair<Fq::Private::PingTaskRequest, double> GetPingTaskRequestWithStatistic(std::optional<FederatedQuery::QueryMeta::ComputeStatus> computeStatus, std::optional<NYql::NDqProto::StatusIds::StatusCode> pendingStatusCode) { - Fq::Private::PingTaskRequest pingTaskRequest; - double cpuUsage = 0.0; - try { - pingTaskRequest = GetPingTaskRequestStatistics(computeStatus, pendingStatusCode, Issues, QueryStats, &cpuUsage); - } catch(const NJson::TJsonException& ex) { - LOG_E("Error statistics conversion: " << ex.what()); - } - - return { pingTaskRequest, cpuUsage }; + void OnPingRequestStart() { + StartTime = TInstant::Now(); + auto pingCounters = Counters.GetCounters(ERequestType::RT_PING); + pingCounters->InFly->Inc(); } void UpdateProgress() { OnPingRequestStart(); - Fq::Private::PingTaskRequest pingTaskRequest = GetPingTaskRequestWithStatistic(std::nullopt, std::nullopt).first; + Fq::Private::PingTaskRequest pingTaskRequest = Builder.Build(QueryStats, Issues); + if (Builder.Issues) { + LOG_W(Builder.Issues.ToOneLineString()); + } + ReportPublicCounters(Builder.PublicStat); Send(Pinger, new TEvents::TEvForwardPingRequest(pingTaskRequest), 0, 1); } @@ -240,8 +242,12 @@ public: LOG_I("Execution status: Failed, Status: " << Status << ", StatusCode: " << NYql::NDqProto::StatusIds::StatusCode_Name(StatusCode) << " Issues: " << Issues.ToOneLineString()); OnPingRequestStart(); - auto [pingTaskRequest, cpuUsage] = GetPingTaskRequestWithStatistic(std::nullopt, StatusCode); - UpdateCpuQuota(cpuUsage); + Fq::Private::PingTaskRequest pingTaskRequest = Builder.Build(QueryStats, Issues, std::nullopt, StatusCode); + if (Builder.Issues) { + LOG_W(Builder.Issues.ToOneLineString()); + } + ReportPublicCounters(Builder.PublicStat); + UpdateCpuQuota(Builder.CpuUsage); Send(Pinger, new TEvents::TEvForwardPingRequest(pingTaskRequest)); } @@ -251,8 +257,12 @@ public: OnPingRequestStart(); ComputeStatus = ::FederatedQuery::QueryMeta::COMPLETING; - auto [pingTaskRequest, cpuUsage] = GetPingTaskRequestWithStatistic(ComputeStatus, std::nullopt); - UpdateCpuQuota(cpuUsage); + Fq::Private::PingTaskRequest pingTaskRequest = Builder.Build(QueryStats, Issues, ComputeStatus, std::nullopt); + if (Builder.Issues) { + LOG_W(Builder.Issues.ToOneLineString()); + } + ReportPublicCounters(Builder.PublicStat); + UpdateCpuQuota(Builder.CpuUsage); Send(Pinger, new TEvents::TEvForwardPingRequest(pingTaskRequest)); } @@ -263,6 +273,7 @@ private: TActorId Connector; TActorId Pinger; NYdb::TOperation::TOperationId OperationId; + PingTaskRequestBuilder Builder; TCounters Counters; NYql::TIssues Issues; NYdb::EStatus Status = NYdb::EStatus::SUCCESS; @@ -271,6 +282,7 @@ private: Ydb::TableStats::QueryStats QueryStats; NKikimr::TBackoffTimer BackoffTimer; FederatedQuery::QueryMeta::ComputeStatus ComputeStatus = FederatedQuery::QueryMeta::RUNNING; + TInstant StartTime; }; std::unique_ptr<NActors::IActor> CreateStatusTrackerActor(const TRunActorParams& params, @@ -278,8 +290,9 @@ std::unique_ptr<NActors::IActor> CreateStatusTrackerActor(const TRunActorParams& const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, + std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters) { - return std::make_unique<TStatusTrackerActor>(params, parent, connector, pinger, operationId, queryCounters); + return std::make_unique<TStatusTrackerActor>(params, parent, connector, pinger, operationId, std::move(processor), queryCounters); } } diff --git a/ydb/core/fq/libs/compute/ydb/status_tracker_actor.h b/ydb/core/fq/libs/compute/ydb/status_tracker_actor.h index a453e2d4d34..f9fc469202c 100644 --- a/ydb/core/fq/libs/compute/ydb/status_tracker_actor.h +++ b/ydb/core/fq/libs/compute/ydb/status_tracker_actor.h @@ -1,6 +1,7 @@ #pragma once #include <ydb/core/fq/libs/compute/common/run_actor_params.h> +#include <ydb/core/fq/libs/compute/common/utils.h> #include <ydb/library/yql/providers/common/metrics/service_counters.h> @@ -13,6 +14,7 @@ std::unique_ptr<NActors::IActor> CreateStatusTrackerActor(const TRunActorParams& const NActors::TActorId& connector, const NActors::TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, + std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters); } diff --git a/ydb/core/fq/libs/compute/ydb/stopper_actor.cpp b/ydb/core/fq/libs/compute/ydb/stopper_actor.cpp index 41e4e73fd4f..c876bcd4422 100644 --- a/ydb/core/fq/libs/compute/ydb/stopper_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/stopper_actor.cpp @@ -1,10 +1,12 @@ -#include "base_status_updater_actor.h" -#include "resources_cleaner_actor.h" +#include "base_compute_actor.h" +#include "stopper_actor.h" +#include <ydb/core/fq/libs/common/compression.h> #include <ydb/core/fq/libs/common/util.h> #include <ydb/core/fq/libs/compute/common/metrics.h> #include <ydb/core/fq/libs/compute/common/retry_actor.h> #include <ydb/core/fq/libs/compute/common/run_actor_params.h> +#include <ydb/core/fq/libs/compute/common/utils.h> #include <ydb/core/fq/libs/compute/ydb/events/events.h> #include <ydb/core/fq/libs/ydb/ydb.h> #include <ydb/library/services/services.pb.h> @@ -32,8 +34,11 @@ namespace NFq { using namespace NActors; using namespace NFq; -class TStopperActor : public TBaseStatusUpdaterActor<TStopperActor> { +class TStopperActor : public TBaseComputeActor<TStopperActor> { public: + + using TBase = TBaseComputeActor<TStopperActor>; + enum ERequestType { RT_CANCEL_OPERATION, RT_GET_OPERATION, @@ -64,17 +69,16 @@ public: } }; - TStopperActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, const ::NYql::NCommon::TServiceCounters& queryCounters) - : TBaseStatusUpdaterActor(params.Config.GetCommon(), queryCounters, "Stopper") + TStopperActor(const TRunActorParams& params, const TActorId& parent, const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters) + : TBase(queryCounters, "Stopper") , Params(params) , Parent(parent) , Connector(connector) , Pinger(pinger) , OperationId(operationId) + , Builder(params.Config.GetCommon(), std::move(processor)) , Counters(GetStepCountersSubgroup()) - { - SetPingCounters(Counters.GetCounters(ERequestType::RT_PING)); - } + {} static constexpr char ActorName[] = "FQ_STOPPER_ACTOR"; @@ -125,20 +129,29 @@ public: auto statusCode = NYql::NDq::YdbStatusToDqStatus(response.StatusCode); LOG_I("Operation successfully fetched, Status: " << response.Status << ", StatusCode: " << NYql::NDqProto::StatusIds::StatusCode_Name(statusCode) << " Issues: " << response.Issues.ToOneLineString()); - OnPingRequestStart(); - Fq::Private::PingTaskRequest pingTaskRequest = GetPingTaskRequest(FederatedQuery::QueryMeta::ABORTING_BY_USER, statusCode, response.Issues, response.QueryStats); + StartTime = TInstant::Now(); + auto pingCounters = Counters.GetCounters(ERequestType::RT_PING); + pingCounters->InFly->Inc(); + + Fq::Private::PingTaskRequest pingTaskRequest = Builder.Build(response.QueryStats, response.Issues, FederatedQuery::QueryMeta::ABORTING_BY_USER, statusCode); + if (Builder.Issues) { + LOG_W(Builder.Issues.ToOneLineString()); + } Send(Pinger, new TEvents::TEvForwardPingRequest(pingTaskRequest)); } void Handle(const TEvents::TEvForwardPingResponse::TPtr& ev) { - OnPingRequestFinish(ev.Get()->Get()->Success); + auto pingCounters = Counters.GetCounters(ERequestType::RT_PING); + pingCounters->InFly->Dec(); + pingCounters->LatencyMs->Collect((TInstant::Now() - StartTime).MilliSeconds()); if (ev.Get()->Get()->Success) { + pingCounters->Ok->Inc(); LOG_I("Information about the status of operation is updated"); } else { + pingCounters->Error->Inc(); LOG_E("Error updating information about the status of operation"); } - Complete(); } @@ -158,7 +171,9 @@ private: TActorId Connector; TActorId Pinger; NYdb::TOperation::TOperationId OperationId; + PingTaskRequestBuilder Builder; TCounters Counters; + TInstant StartTime; }; std::unique_ptr<NActors::IActor> CreateStopperActor(const TRunActorParams& params, @@ -166,8 +181,9 @@ std::unique_ptr<NActors::IActor> CreateStopperActor(const TRunActorParams& param const TActorId& connector, const TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, + std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters) { - return std::make_unique<TStopperActor>(params, parent, connector, pinger, operationId, queryCounters); + return std::make_unique<TStopperActor>(params, parent, connector, pinger, operationId, std::move(processor), queryCounters); } } diff --git a/ydb/core/fq/libs/compute/ydb/stopper_actor.h b/ydb/core/fq/libs/compute/ydb/stopper_actor.h index 4f41a3e5dcf..f078664566c 100644 --- a/ydb/core/fq/libs/compute/ydb/stopper_actor.h +++ b/ydb/core/fq/libs/compute/ydb/stopper_actor.h @@ -1,6 +1,7 @@ #pragma once #include <ydb/core/fq/libs/compute/common/run_actor_params.h> +#include <ydb/core/fq/libs/compute/common/utils.h> #include <ydb/library/yql/providers/common/metrics/service_counters.h> @@ -13,6 +14,7 @@ std::unique_ptr<NActors::IActor> CreateStopperActor(const TRunActorParams& param const NActors::TActorId& connector, const NActors::TActorId& pinger, const NYdb::TOperation::TOperationId& operationId, + std::unique_ptr<IPlanStatProcessor>&& processor, const ::NYql::NCommon::TServiceCounters& queryCounters); } diff --git a/ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp b/ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp index 2b52e783378..e0d7c520014 100644 --- a/ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp +++ b/ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp @@ -24,12 +24,7 @@ public: : YqSharedResources(params.YqSharedResources) , CredentialsProviderFactory(params.CredentialsProviderFactory) , ComputeConnection(params.ComputeConnection) - { - StatsMode = params.Config.GetControlPlaneStorage().GetStatsMode(); - if (StatsMode == Ydb::Query::StatsMode::STATS_MODE_UNSPECIFIED) { - StatsMode = Ydb::Query::StatsMode::STATS_MODE_FULL; - } - } + {} void Bootstrap() { auto querySettings = NFq::GetClientSettings<NYdb::NQuery::TClientSettings>(ComputeConnection, CredentialsProviderFactory); @@ -55,7 +50,7 @@ public: settings.OperationTimeout(event.OperationTimeout); settings.Syntax(event.Syntax); settings.ExecMode(event.ExecMode); - settings.StatsMode(StatsMode); + settings.StatsMode(event.StatsMode); settings.TraceId(event.TraceId); QueryClient ->ExecuteScript(event.Sql, settings) @@ -222,7 +217,6 @@ private: NConfig::TYdbStorageConfig ComputeConnection; std::unique_ptr<NYdb::NQuery::TQueryClient> QueryClient; std::unique_ptr<NYdb::NOperation::TOperationClient> OperationClient; - Ydb::Query::StatsMode StatsMode; }; std::unique_ptr<NActors::IActor> CreateConnectorActor(const TRunActorParams& params) { diff --git a/ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp b/ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp index 246d3f3852d..c0802446d01 100644 --- a/ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp +++ b/ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp @@ -256,7 +256,8 @@ TPingTaskParams ConstructHardPingTask( internal.clear_statistics(); PackStatisticsToProtobuf(*internal.mutable_statistics(), statistics); - if (!dumpRawStatistics) { + // global dumpRawStatistics will be removed with YQv1 + if (!dumpRawStatistics && !request.dump_raw_statistics()) { try { statistics = GetPrettyStatistics(statistics); } catch (const std::exception&) { diff --git a/ydb/core/fq/libs/protos/fq_private.proto b/ydb/core/fq/libs/protos/fq_private.proto index bba6f0ffdb4..61d9f492565 100644 --- a/ydb/core/fq/libs/protos/fq_private.proto +++ b/ydb/core/fq/libs/protos/fq_private.proto @@ -162,6 +162,7 @@ message PingTaskRequest { string operation_id = 35; string execution_id = 36; NYql.NDqProto.StatusIds.StatusCode pending_status_code = 37; + bool dump_raw_statistics = 38; } message PingTaskResult { |
