summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorVitalii Gridnev <[email protected]>2024-07-26 12:58:09 +0300
committerGitHub <[email protected]>2024-07-26 12:58:09 +0300
commit7cf9db705d539dbf732ea7e8d35e8a8a2579fff1 (patch)
tree7ec4c8b8380dea3989a99acdb009c14a51a5804d
parent6064125b9de292c07d79f1613b49985cde04d682 (diff)
add feature flag to enable spilling nodes (#6895)
-rw-r--r--ydb/core/kqp/compile_service/kqp_compile_actor.cpp1
-rw-r--r--ydb/core/kqp/compile_service/kqp_compile_service.cpp3
-rw-r--r--ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp4
-rw-r--r--ydb/core/kqp/executer_actor/kqp_planner.cpp1
-rw-r--r--ydb/core/kqp/provider/yql_kikimr_settings.cpp44
-rw-r--r--ydb/core/kqp/provider/yql_kikimr_settings.h5
-rw-r--r--ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp60
-rw-r--r--ydb/core/protos/table_service_config.proto2
8 files changed, 100 insertions, 20 deletions
diff --git a/ydb/core/kqp/compile_service/kqp_compile_actor.cpp b/ydb/core/kqp/compile_service/kqp_compile_actor.cpp
index aa2427615cf..679120cfe09 100644
--- a/ydb/core/kqp/compile_service/kqp_compile_actor.cpp
+++ b/ydb/core/kqp/compile_service/kqp_compile_actor.cpp
@@ -608,6 +608,7 @@ void ApplyServiceConfig(TKikimrConfiguration& kqpConfig, const TTableServiceConf
kqpConfig.EnableSpillingGenericQuery = serviceConfig.GetEnableQueryServiceSpilling();
kqpConfig.DefaultCostBasedOptimizationLevel = serviceConfig.GetDefaultCostBasedOptimizationLevel();
kqpConfig.EnableConstantFolding = serviceConfig.GetEnableConstantFolding();
+ kqpConfig.SetDefaultEnabledSpillingNodes(serviceConfig.GetEnableSpillingNodes());
if (const auto limit = serviceConfig.GetResourceManager().GetMkqlHeavyProgramMemoryLimit()) {
kqpConfig._KqpYqlCombinerMemoryLimit = std::max(1_GB, limit - (limit >> 2U));
diff --git a/ydb/core/kqp/compile_service/kqp_compile_service.cpp b/ydb/core/kqp/compile_service/kqp_compile_service.cpp
index 0fac7824ef5..d0788088ff0 100644
--- a/ydb/core/kqp/compile_service/kqp_compile_service.cpp
+++ b/ydb/core/kqp/compile_service/kqp_compile_service.cpp
@@ -534,6 +534,8 @@ private:
ui64 defaultCostBasedOptimizationLevel = TableServiceConfig.GetDefaultCostBasedOptimizationLevel();
bool enableConstantFolding = TableServiceConfig.GetEnableConstantFolding();
+ TString enableSpillingNodes = TableServiceConfig.GetEnableSpillingNodes();
+
TableServiceConfig.Swap(event.MutableConfig()->MutableTableServiceConfig());
LOG_INFO(*TlsActivationContext, NKikimrServices::KQP_COMPILE_SERVICE, "Updated config");
@@ -556,6 +558,7 @@ private:
TableServiceConfig.GetExtractPredicateRangesLimit() != rangesLimit ||
TableServiceConfig.GetResourceManager().GetMkqlHeavyProgramMemoryLimit() != mkqlHeavyLimit ||
TableServiceConfig.GetIdxLookupJoinPointsLimit() != idxLookupPointsLimit ||
+ TableServiceConfig.GetEnableSpillingNodes() != enableSpillingNodes ||
TableServiceConfig.GetEnableQueryServiceSpilling() != enableQueryServiceSpilling ||
TableServiceConfig.GetDefaultCostBasedOptimizationLevel() != defaultCostBasedOptimizationLevel ||
TableServiceConfig.GetEnableConstantFolding() != enableConstantFolding ||
diff --git a/ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp b/ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp
index ca920e7112d..c97a914f0c0 100644
--- a/ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp
+++ b/ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp
@@ -166,6 +166,10 @@ public:
runtimeSettings.UseSpilling = args.WithSpilling;
runtimeSettings.StatsMode = args.StatsMode;
+ if (runtimeSettings.UseSpilling) {
+ args.Task->SetEnableSpilling(runtimeSettings.UseSpilling);
+ }
+
if (args.Deadline) {
runtimeSettings.Timeout = args.Deadline - TAppData::TimeProvider->Now();
}
diff --git a/ydb/core/kqp/executer_actor/kqp_planner.cpp b/ydb/core/kqp/executer_actor/kqp_planner.cpp
index 7065b49936f..292d05c65dc 100644
--- a/ydb/core/kqp/executer_actor/kqp_planner.cpp
+++ b/ydb/core/kqp/executer_actor/kqp_planner.cpp
@@ -194,6 +194,7 @@ std::unique_ptr<TEvKqpNode::TEvStartKqpTasksRequest> TKqpPlanner::SerializeReque
request.SetStartAllOrFail(true);
if (UseDataQueryPool) {
request.MutableRuntimeSettings()->SetExecType(NYql::NDqProto::TComputeRuntimeSettings::DATA);
+ request.MutableRuntimeSettings()->SetUseSpilling(WithSpilling);
} else {
request.MutableRuntimeSettings()->SetExecType(NYql::NDqProto::TComputeRuntimeSettings::SCAN);
request.MutableRuntimeSettings()->SetUseSpilling(WithSpilling);
diff --git a/ydb/core/kqp/provider/yql_kikimr_settings.cpp b/ydb/core/kqp/provider/yql_kikimr_settings.cpp
index 9ccf65525dd..960e464e1aa 100644
--- a/ydb/core/kqp/provider/yql_kikimr_settings.cpp
+++ b/ydb/core/kqp/provider/yql_kikimr_settings.cpp
@@ -25,6 +25,22 @@ EOptionalFlag GetOptionalFlagValue(const TMaybe<TType>& flag) {
return EOptionalFlag::Disabled;
}
+
+ui64 ParseEnableSpillingNodes(const TString &v) {
+ ui64 res = 0;
+ TVector<TString> vec;
+ StringSplitter(v).SplitBySet(",;| ").AddTo(&vec);
+ for (auto& s: vec) {
+ if (s.empty()) {
+ throw yexception() << "Empty value item";
+ }
+ auto value = FromStringWithDefault<NYql::TDqSettings::EEnabledSpillingNodes>(
+ s, NYql::TDqSettings::EEnabledSpillingNodes::None);
+ res |= ui64(value);
+ }
+ return res;
+}
+
static inline bool GetFlagValue(const TMaybe<bool>& flag) {
return flag ? flag.GetRef() : false;
}
@@ -73,20 +89,7 @@ TKikimrConfiguration::TKikimrConfiguration() {
REGISTER_SETTING(*this, OptUseFinalizeByKey);
REGISTER_SETTING(*this, CostBasedOptimizationLevel);
REGISTER_SETTING(*this, EnableSpillingNodes)
- .Parser([](const TString& v) {
- ui64 res = 0;
- TVector<TString> vec;
- StringSplitter(v).SplitBySet(",;| ").AddTo(&vec);
- for (auto& s: vec) {
- if (s.empty()) {
- throw yexception() << "Empty value item";
- }
- auto value = FromStringWithDefault<NYql::TDqSettings::EEnabledSpillingNodes>(
- s, NYql::TDqSettings::EEnabledSpillingNodes::None);
- res |= ui64(value);
- }
- return res;
- });
+ .Parser([](const TString& v) { return ParseEnableSpillingNodes(v); });
REGISTER_SETTING(*this, MaxDPccpDPTableSize);
@@ -143,11 +146,6 @@ bool TKikimrSettings::HasOptUseFinalizeByKey() const {
return GetOptionalFlagValue(OptUseFinalizeByKey.Get()) != EOptionalFlag::Disabled;
}
-ui64 TKikimrSettings::GetEnabledSpillingNodes() const {
- return EnableSpillingNodes.Get().GetOrElse(0);
-}
-
-
EOptionalFlag TKikimrSettings::GetOptPredicateExtract() const {
return GetOptionalFlagValue(OptEnablePredicateExtract.Get());
}
@@ -169,4 +167,12 @@ TKikimrSettings::TConstPtr TKikimrConfiguration::Snapshot() const {
return std::make_shared<const TKikimrSettings>(*this);
}
+void TKikimrConfiguration::SetDefaultEnabledSpillingNodes(const TString& node) {
+ DefaultEnableSpillingNodes = ParseEnableSpillingNodes(node);
+}
+
+ui64 TKikimrConfiguration::GetEnabledSpillingNodes() const {
+ return EnableSpillingNodes.Get().GetOrElse(DefaultEnableSpillingNodes);
+}
+
}
diff --git a/ydb/core/kqp/provider/yql_kikimr_settings.h b/ydb/core/kqp/provider/yql_kikimr_settings.h
index 19699a317d5..b84df708900 100644
--- a/ydb/core/kqp/provider/yql_kikimr_settings.h
+++ b/ydb/core/kqp/provider/yql_kikimr_settings.h
@@ -84,7 +84,6 @@ struct TKikimrSettings {
bool HasOptEnableOlapPushdown() const;
bool HasOptEnableOlapProvideComputeSharding() const;
bool HasOptUseFinalizeByKey() const;
- ui64 GetEnabledSpillingNodes() const;
EOptionalFlag GetOptPredicateExtract() const;
EOptionalFlag GetUseLlvm() const;
@@ -167,6 +166,10 @@ struct TKikimrConfiguration : public TKikimrSettings, public NCommon::TSettingDi
bool EnableSpillingGenericQuery = false;
ui32 DefaultCostBasedOptimizationLevel = 3;
bool EnableConstantFolding = true;
+ ui64 DefaultEnableSpillingNodes = 0;
+
+ void SetDefaultEnabledSpillingNodes(const TString& node);
+ ui64 GetEnabledSpillingNodes() const;
};
}
diff --git a/ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp b/ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp
index 91c0f3b7610..c092ca9ef6c 100644
--- a/ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp
+++ b/ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp
@@ -32,10 +32,70 @@ NKikimrConfig::TAppConfig AppCfg() {
return appCfg;
}
+NKikimrConfig::TAppConfig AppCfgLowComputeLimits(ui64 reasonableTreshold) {
+ NKikimrConfig::TAppConfig appCfg;
+
+ auto* rm = appCfg.MutableTableServiceConfig()->MutableResourceManager();
+ rm->SetMkqlLightProgramMemoryLimit(100);
+ rm->SetMkqlHeavyProgramMemoryLimit(300);
+ rm->SetReasonableSpillingTreshold(reasonableTreshold);
+ appCfg.MutableTableServiceConfig()->SetEnableQueryServiceSpilling(true);
+
+ auto* spilling = appCfg.MutableTableServiceConfig()->MutableSpillingServiceConfig()->MutableLocalFileConfig();
+
+ spilling->SetEnable(true);
+ spilling->SetRoot("./spilling/");
+
+ return appCfg;
+}
+
+
} // anonymous namespace
Y_UNIT_TEST_SUITE(KqpScanSpilling) {
+Y_UNIT_TEST_TWIN(SpillingInRuntimeNodes, EnabledSpilling) {
+ ui64 reasonableTreshold = EnabledSpilling ? 100 : 200_MB;
+ Cerr << "cwd: " << NFs::CurrentWorkingDirectory() << Endl;
+ TKikimrRunner kikimr(AppCfgLowComputeLimits(reasonableTreshold));
+
+ auto db = kikimr.GetQueryClient();
+
+ for (ui32 i = 0; i < 300; ++i) {
+ auto result = db.ExecuteQuery(Sprintf(R"(
+ --!syntax_v1
+ REPLACE INTO `/Root/KeyValue` (Key, Value) VALUES (%d, "%s")
+ )", i, TString(200000 + i, 'a' + (i % 26)).c_str()), NYdb::NQuery::TTxControl::BeginTx().CommitTx()).GetValueSync();
+ UNIT_ASSERT_C(result.IsSuccess(), result.GetIssues().ToString());
+ }
+
+ auto query = R"(
+ --!syntax_v1
+ PRAGMA ydb.EnableSpillingNodes="GraceJoin";
+ select t1.Key, t1.Value, t2.Key, t2.Value
+ from `/Root/KeyValue` as t1 full join `/Root/KeyValue` as t2 on t1.Value = t2.Value
+ order by t1.Value
+ )";
+
+ auto explainMode = NYdb::NQuery::TExecuteQuerySettings().ExecMode(NYdb::NQuery::EExecMode::Explain);
+ auto planres = db.ExecuteQuery(query, NYdb::NQuery::TTxControl::NoTx(), explainMode).ExtractValueSync();
+ UNIT_ASSERT_VALUES_EQUAL_C(planres.GetStatus(), EStatus::SUCCESS, planres.GetIssues().ToString());
+
+ Cerr << planres.GetStats()->GetAst() << Endl;
+
+ auto result = db.ExecuteQuery(query, NYdb::NQuery::TTxControl::BeginTx().CommitTx(), NYdb::NQuery::TExecuteQuerySettings()).ExtractValueSync();
+ UNIT_ASSERT_VALUES_EQUAL_C(result.GetStatus(), EStatus::SUCCESS, result.GetIssues().ToString());
+
+ TKqpCounters counters(kikimr.GetTestServer().GetRuntime()->GetAppData().Counters);
+ if (EnabledSpilling) {
+ UNIT_ASSERT(counters.SpillingWriteBlobs->Val() > 0);
+ UNIT_ASSERT(counters.SpillingReadBlobs->Val() > 0);
+ } else {
+ UNIT_ASSERT(counters.SpillingWriteBlobs->Val() == 0);
+ UNIT_ASSERT(counters.SpillingReadBlobs->Val() == 0);
+ }
+}
+
Y_UNIT_TEST(SelfJoinQueryService) {
Cerr << "cwd: " << NFs::CurrentWorkingDirectory() << Endl;
diff --git a/ydb/core/protos/table_service_config.proto b/ydb/core/protos/table_service_config.proto
index be6e5b362de..0ee2c712f9c 100644
--- a/ydb/core/protos/table_service_config.proto
+++ b/ydb/core/protos/table_service_config.proto
@@ -300,4 +300,6 @@ message TTableServiceConfig {
optional bool EnableConstantFolding = 65 [ default = true ];
optional bool EnableImplicitQueryParameterTypes = 66 [ default = true ];
+
+ optional string EnableSpillingNodes = 67 [ default = "All" ];
};