summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
-rw-r--r--ydb/core/fq/libs/compute/common/utils.cpp144
-rw-r--r--ydb/core/fq/libs/compute/common/utils.h33
-rw-r--r--ydb/core/fq/libs/compute/ydb/actors_factory.cpp44
-rw-r--r--ydb/core/fq/libs/compute/ydb/base_status_updater_actor.h96
-rw-r--r--ydb/core/fq/libs/compute/ydb/events/events.h4
-rw-r--r--ydb/core/fq/libs/compute/ydb/executer_actor.cpp9
-rw-r--r--ydb/core/fq/libs/compute/ydb/executer_actor.h1
-rw-r--r--ydb/core/fq/libs/compute/ydb/status_tracker_actor.cpp131
-rw-r--r--ydb/core/fq/libs/compute/ydb/status_tracker_actor.h2
-rw-r--r--ydb/core/fq/libs/compute/ydb/stopper_actor.cpp42
-rw-r--r--ydb/core/fq/libs/compute/ydb/stopper_actor.h2
-rw-r--r--ydb/core/fq/libs/compute/ydb/ydb_connector_actor.cpp10
-rw-r--r--ydb/core/fq/libs/control_plane_storage/internal/task_ping.cpp3
-rw-r--r--ydb/core/fq/libs/protos/fq_private.proto1
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 {