diff options
| author | Vitalii Gridnev <[email protected]> | 2024-07-26 12:58:09 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-07-26 12:58:09 +0300 |
| commit | 7cf9db705d539dbf732ea7e8d35e8a8a2579fff1 (patch) | |
| tree | 7ec4c8b8380dea3989a99acdb009c14a51a5804d | |
| parent | 6064125b9de292c07d79f1613b49985cde04d682 (diff) | |
add feature flag to enable spilling nodes (#6895)
| -rw-r--r-- | ydb/core/kqp/compile_service/kqp_compile_actor.cpp | 1 | ||||
| -rw-r--r-- | ydb/core/kqp/compile_service/kqp_compile_service.cpp | 3 | ||||
| -rw-r--r-- | ydb/core/kqp/compute_actor/kqp_compute_actor_factory.cpp | 4 | ||||
| -rw-r--r-- | ydb/core/kqp/executer_actor/kqp_planner.cpp | 1 | ||||
| -rw-r--r-- | ydb/core/kqp/provider/yql_kikimr_settings.cpp | 44 | ||||
| -rw-r--r-- | ydb/core/kqp/provider/yql_kikimr_settings.h | 5 | ||||
| -rw-r--r-- | ydb/core/kqp/ut/spilling/kqp_scan_spilling_ut.cpp | 60 | ||||
| -rw-r--r-- | ydb/core/protos/table_service_config.proto | 2 |
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" ]; }; |
