diff options
| author | Sergey Uzhakov <[email protected]> | 2026-07-15 14:54:14 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-15 14:54:14 +0300 |
| commit | 1dd74a25bf627151fced1f37a0045205b9e48405 (patch) | |
| tree | 8a0efe433e6d83a0bb095470bf25b7852fb9f38c | |
| parent | cf6219c66d10cb909b70c559cf9de6a115cf3fab (diff) | |
remove dqrun and unused dq-related code (#46501)
46 files changed, 36 insertions, 5989 deletions
diff --git a/ydb/library/yql/providers/dq/actors/ya.make b/ydb/library/yql/providers/dq/actors/ya.make index 3988a3b10f3..f4aa53c3872 100644 --- a/ydb/library/yql/providers/dq/actors/ya.make +++ b/ydb/library/yql/providers/dq/actors/ya.make @@ -29,7 +29,6 @@ PEERDIR( ydb/library/yql/providers/dq/api/grpc ydb/library/yql/providers/dq/api/protos ydb/library/yql/providers/dq/common - ydb/library/yql/providers/dq/config ydb/library/yql/providers/dq/counters ydb/library/yql/providers/dq/interface ydb/library/yql/providers/dq/planner diff --git a/ydb/library/yql/providers/dq/common/yql_dq_common.h b/ydb/library/yql/providers/dq/common/yql_dq_common.h index d9601d233db..abb2da89618 100644 --- a/ydb/library/yql/providers/dq/common/yql_dq_common.h +++ b/ydb/library/yql/providers/dq/common/yql_dq_common.h @@ -8,6 +8,12 @@ #include <map> namespace NYql { + +inline const TString& DqStrippedSuffied() { + static const TString suffix(".s"); + return suffix; +} + namespace NCommon { struct TResultFormatSettings { diff --git a/ydb/library/yql/providers/dq/config/config.proto b/ydb/library/yql/providers/dq/config/config.proto deleted file mode 100644 index b3e86fa40db..00000000000 --- a/ydb/library/yql/providers/dq/config/config.proto +++ /dev/null @@ -1,193 +0,0 @@ -package NYql.NProto; -option java_package = "ru.yandex.yql.proto"; -option go_package = "github.com/ydb-platform/ydb/ydb/library/yql/providers/dq/config;config_pb"; - -message TDqConfig { - optional uint32 NodeId = 1; // autofilled by yqlworker - optional uint32 ActorThreads = 2 [default = 4]; - optional string Host = 3; // autofilled by yqlworker - optional uint32 Port = 4; // autofilled by yqlworker - - message TICSettings { - optional uint64 HandshakeMs = 1 [default = 5000]; - optional uint64 DeadPeerMs = 2 [default = 30000]; - optional uint64 CloseOnIdleMs = 3 [default = 300000]; - optional uint32 SendBufferDieLimitInMB = 4 [default = 512]; - optional uint64 OutputBuffersTotalSizeLimitInMB = 5 [default = 0]; - optional uint32 MaxInflightAmountOfData = 6 [default = 10485760]; - optional uint32 TotalInflightAmountOfData = 7 [default = 0]; - optional bool MergePerPeerCounters = 8 [default = false]; - optional bool MergePerDataCenterCounters = 9 [default = false]; - optional uint32 TCPSocketBufferSize = 10 [default = 16777216]; - optional uint64 PingPeriodMs = 11 [default = 3000]; - optional uint64 ForceConfirmPeriodMs = 12 [default = 1000]; - optional uint64 LostConnectionMs = 13; - optional uint64 BatchPeriodMs = 14; - optional uint64 MessagePendingTimeoutMs = 15 [default = 5000]; - optional uint64 MessagePendingSize = 16 [default = 18446744073709551615]; - optional uint32 MaxSerializedEventSize = 17 [default = 67108000]; - optional bool EnableExternalDataChannel = 27 [default = false]; - - // Scheduler - optional uint64 ResolutionMicroseconds = 18 [default = 1024]; - optional uint64 SpinThreshold = 19 [default = 100]; - optional uint64 ProgressThreshold = 20 [default = 10000]; - optional bool UseSchedulerActor = 21 [default = false]; - optional uint64 RelaxedSendPaceEventsPerSecond = 22; - optional uint64 RelaxedSendPaceEventsPerCycle = 23; - optional uint64 RelaxedSendThresholdEventsPerSecond = 24; - optional uint64 RelaxedSendThresholdEventsPerCycle = 25; - - // Thread Pool - optional uint32 Threads = 26 [default = 4]; - } - - message TSolomon { - optional string TokenFile = 1; - optional string Token = 2; - optional string Server = 3 [default = "https://solomon.yandex.net"]; - optional string Path = 4 [default = "/api/v2/push"]; - optional int32 Port = 5 [default = 443]; - optional string Project = 6 [default = "yql"]; - optional string Service = 7 [default = "dq_vanilla"]; - optional string Cluster = 8 [default = "test"]; - } - - message TAttr { - optional string Name = 1; - optional string Value = 2; - } - - message TFile { - optional string Name = 1; - optional string LocalPath = 2; - } - - message TPortoSettings { - repeated TAttr Setting = 1; - } - - message TSpillingSettings { - optional string Root = 1 [default = "./spilling"]; - - optional uint64 MaxTotalSize = 2; - optional uint64 MaxFileSize = 3; - optional uint64 MaxFilePartSize = 4; - - optional uint32 IoThreadPoolWorkersCount = 5 [default = 2]; - optional uint32 IoThreadPoolQueueSize = 6 [default = 1000]; - optional bool CleanupOnShutdown = 7; - } - - message TDiskRequest { - optional uint64 DiskSpace = 1; - optional uint64 InodeCount = 2; - optional string Account = 3; - optional string MediumName = 4; - optional string AdditionalSpecYson = 5; // YSON text map used as disk_request base; structured fields overwrite matching keys (e.g. add nbd_disk in YSON, set medium_name in DiskRequest) - } - - message TYtBackend { - optional string ClusterName = 1 [default = "hume"]; - optional string User = 2; // default -- current user name - optional string TokenFile = 3; // default -- $HOME/.yt/token - optional string VanillaJobLite = 4; - optional string VanillaJobLiteMd5 = 27; // autofilled (for Lite version only) - optional string VanillaJobCommand = 35; - repeated TFile VanillaJobFile = 36; - optional uint32 MaxJobs = 5; - optional string Prefix = 16; - optional string UploadPrefix = 6; // deprecated option - optional string Pool = 7; - repeated string PoolTrees = 24; - optional string Token = 8; // for internal use only - optional int64 MemoryLimit = 9; - optional int64 CpuLimit = 29; - repeated TAttr VaultEnv = 34; - optional uint32 JobsPerOperation = 11; - optional uint32 UploadReplicationFactor = 12; - repeated string Owner = 13; - optional uint32 MinNodeId = 14; // inclusive - optional uint32 MaxNodeId = 15; // exclusive - optional uint32 ActorStartPort = 17 [default = 31002]; - optional bool SameActorPorts = 18 [default = true]; - optional bool UseTmpFs = 19; - optional int64 CacheSize = 20; - optional string NetworkProject = 21; - optional int32 WorkerCapacity = 22 [default = 1]; - optional TICSettings ICSettings = 23; // can be filled by yqlworker - optional string EnablePorto = 25; - repeated string PortoLayer = 28; - repeated string OperationLayer = 37; - optional TPortoSettings PortoSettings = 31; - optional bool ContainerCpuLimit = 26; // for testing only - optional TSolomon Solomon = 30; - optional bool CanUseComputeActor = 32 [default = false]; - optional bool EnforceJobUtc = 33; - optional TSpillingSettings SpillingSettings = 38; - optional TDiskRequest DiskRequest = 39; - optional bool EnforceJobYtIsolation = 40; - optional bool UseLocalLDLibraryPath = 41 [default = false]; - optional string SchedulingTagFilter = 42; - optional uint64 EnforceRegexpProbabilityFail = 43; // Same as _EnforceRegexpProbabilityFail. - optional string ProxyAddress = 44; // full URL to YT proxy, e.g. https://my.com - } - - repeated TYtBackend YtBackends = 5; - - message TYtCoordinator { - optional string ClusterName = 1 [default = "hume"]; - optional string User = 2; // default -- current user name - optional string TokenFile = 3; // default -- $HOME/.yt/token - optional string Token = 4; // for internal use only - optional string Prefix = 5; - optional string DebugLogFile = 6; - optional string HostName = 7; // for debug only - optional int64 HeartbeatPeriodMs = 8 [default = 2000]; - repeated string ServiceNodeHostPort = 9; // for tests and debug only - optional string Revision = 10; // for tests and debug only - optional string LockType = 11 [default = "yt"]; // for tests and debug only - optional string ProxyAddress = 12; // full URL to YT proxy, e.g. https://my.com - } - - optional TYtCoordinator YtCoordinator = 6; - - optional uint32 PortStart = 7; - optional uint32 PortFinish = 8; - - optional bool DisableNodeCleaner = 9; // for tests - - message TScheduler { - optional bool KeepReserveForLiteralRequests = 1 [default = true]; - optional uint32 HistoryKeepingTime = 2 [default = 60]; // in minutes - optional uint32 MaxOperations = 3 [default = 1000]; - optional uint32 MaxOperationsPerUser = 4 [default = 100]; - optional bool LimitTasksPerWindow = 5 [default = false]; - optional uint32 LimiterNumerator = 6 [default = 1]; - optional uint32 LimiterDenumerator = 7 [default = 2]; - optional uint32 MaxRequestsPerTick = 8 [default = 200]; - } - - optional TScheduler Scheduler = 10; - - message TDqControl { - optional bool Disabled = 1 [default = false]; - repeated string IndexedUdfsToWarmup = 2; - optional bool WaitOnStart = 3 [default = false]; - optional bool EnableStrip = 4 [default = true]; - optional bool EnableCliqueWarmup = 5 [default = false]; - } - - optional TDqControl Control = 11; - - optional TICSettings ICSettings = 12; - - optional string User = 13; - optional string Group = 14; - - optional TSolomon Solomon = 15; - optional TPortoSettings PortoSettings = 16; - - optional uint64 OpenSessionTimeoutMs = 17 [default = 15000]; - optional uint64 RequestTimeoutMs = 18 [default = 720000]; -} diff --git a/ydb/library/yql/providers/dq/config/ya.make b/ydb/library/yql/providers/dq/config/ya.make deleted file mode 100644 index d302ee72c99..00000000000 --- a/ydb/library/yql/providers/dq/config/ya.make +++ /dev/null @@ -1,8 +0,0 @@ -PROTO_LIBRARY() -PROTOC_FATAL_WARNINGS() - -SRCS( - config.proto -) - -END() diff --git a/ydb/library/yql/providers/dq/local_gateway/ut/ya.make b/ydb/library/yql/providers/dq/local_gateway/ut/ya.make deleted file mode 100644 index a58c5a8f789..00000000000 --- a/ydb/library/yql/providers/dq/local_gateway/ut/ya.make +++ /dev/null @@ -1,23 +0,0 @@ -UNITTEST_FOR(ydb/library/yql/providers/dq/local_gateway) - -SRCS( - yql_dq_gateway_local_ut.cpp -) - -PEERDIR( - yql/essentials/public/udf/service/stub - yql/essentials/sql/pg_dummy - ydb/library/yql/dq/comp_nodes/no_llvm - ydb/library/yql/dq/runtime - ydb/library/yql/providers/dq/local_gateway - ydb/library/yql/dq/transform - yql/essentials/providers/common/comp_nodes - yql/essentials/minikql - yql/essentials/minikql/invoke_builtins/no_llvm - yql/essentials/minikql/comp_nodes/no_llvm - yql/essentials/minikql/computation/no_llvm -) - -YQL_LAST_ABI_VERSION() - -END() diff --git a/ydb/library/yql/providers/dq/local_gateway/ut/yql_dq_gateway_local_ut.cpp b/ydb/library/yql/providers/dq/local_gateway/ut/yql_dq_gateway_local_ut.cpp deleted file mode 100644 index 00f8052723a..00000000000 --- a/ydb/library/yql/providers/dq/local_gateway/ut/yql_dq_gateway_local_ut.cpp +++ /dev/null @@ -1,56 +0,0 @@ -#include <library/cpp/testing/unittest/registar.h> - -#include <ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h> -#include <ydb/library/yql/dq/transform/yql_common_dq_transform.h> -#include <yql/essentials/minikql/mkql_function_registry.h> -#include <yql/essentials/minikql/invoke_builtins/mkql_builtins.h> -#include <ydb/library/yql/dq/comp_nodes/yql_common_dq_factory.h> -#include <yql/essentials/providers/common/comp_nodes/yql_factory.h> -#include <yql/essentials/minikql/comp_nodes/mkql_factories.h> -#include <yql/essentials/minikql/mkql_function_registry.h> - -using namespace NYql; -using namespace NKikimr; - -Y_UNIT_TEST_SUITE(TestLocalGateway) { - -std::pair<TIntrusivePtr<IDqGateway>,TIntrusivePtr<NMiniKQL::IFunctionRegistry>> CreateGateway() -{ - auto funcRegistry = CreateFunctionRegistry(NMiniKQL::CreateBuiltinRegistry()); - - auto dqTaskTransformFactory = CreateCommonDqTaskTransformFactory(); - - auto dqCompFactory = GetCommonDqFactory(); - - TDqTaskPreprocessorFactoryCollection dqTaskPreprocessorFactories; - - Y_UNUSED(dqTaskTransformFactory); - - return {CreateLocalDqGateway( - funcRegistry.Get(), - dqCompFactory, - dqTaskTransformFactory, - dqTaskPreprocessorFactories, - /*enableSpilling = */ false, - MakeIntrusive<NYql::NDq::TDqAsyncIoFactory>()), funcRegistry}; -} - -Y_UNIT_TEST(Create) { - auto [gateway, _] = CreateGateway(); -} - -Y_UNIT_TEST(OpenSession) { - auto [gateway, _] = CreateGateway(); - gateway->OpenSession("SessionId", "Username"); -} - -Y_UNIT_TEST(OpenSessionInLoop) { - for (int i = 0; i < 10; i++) { - Cerr << "Iteration: " << i << "\n"; - auto [gateway, _] = CreateGateway(); - gateway->OpenSession("SessionId", "Username"); - } -} - -} // Y_UNIT_TEST_SUITE(TestLocalGateway) - diff --git a/ydb/library/yql/providers/dq/local_gateway/ya.make b/ydb/library/yql/providers/dq/local_gateway/ya.make deleted file mode 100644 index 116de3e6de7..00000000000 --- a/ydb/library/yql/providers/dq/local_gateway/ya.make +++ /dev/null @@ -1,22 +0,0 @@ -LIBRARY() - -YQL_LAST_ABI_VERSION() - -SRCS( - yql_dq_gateway_local.cpp -) - -PEERDIR( - yql/essentials/utils - yql/essentials/utils/network - ydb/library/yql/dq/actors/compute - ydb/library/yql/dq/actors/spilling - ydb/library/yql/providers/dq/provider - ydb/library/yql/providers/dq/api/protos - ydb/library/yql/providers/dq/task_runner - ydb/library/yql/providers/dq/worker_manager - ydb/library/yql/providers/dq/service - ydb/library/yql/providers/dq/stats_collector -) - -END() diff --git a/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.cpp b/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.cpp deleted file mode 100644 index 5c5693de291..00000000000 --- a/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.cpp +++ /dev/null @@ -1,302 +0,0 @@ -#include "yql_dq_gateway_local.h" - -#include <ydb/library/yql/providers/dq/provider/yql_dq_gateway.h> -#include <ydb/library/yql/providers/dq/task_runner/tasks_runner_local.h> - -#include <ydb/library/yql/providers/dq/service/interconnect_helpers.h> -#include <ydb/library/yql/providers/dq/service/service_node.h> - -#include <ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.h> -#include <ydb/library/yql/providers/dq/worker_manager/local_worker_manager.h> -#include <ydb/library/yql/dq/actors/spilling/spilling_file.h> - -#include <yql/essentials/utils/range_walker.h> -#include <yql/essentials/utils/network/bind_in_range.h> - -#include <library/cpp/messagebus/network.h> - -#include <util/system/env.h> -#include <util/generic/size_literals.h> -#include <util/folder/dirut.h> - -namespace NYql { - -using namespace NActors; -using NDqs::MakeWorkerManagerActorID; - -class TLocalServiceHolder { -public: - TLocalServiceHolder(const NKikimr::NMiniKQL::IFunctionRegistry* functionRegistry, NKikimr::NMiniKQL::TComputationNodeFactory compFactory, - TTaskTransformFactory taskTransformFactory, const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories, NBus::TBindResult interconnectPort, NBus::TBindResult grpcPort, - NDq::IDqAsyncIoFactory::TPtr asyncIoFactory, int threads, - IMetricsRegistryPtr metricsRegistry, - const std::function<IActor*(void)>& metricsPusherFactory, - bool withSpilling, - TVector<std::pair<TActorId, TActorSetupCmd>>&& additionalLocalServices) - : MetricsRegistry(metricsRegistry - ? metricsRegistry - : CreateMetricsRegistry(GetSensorsGroupFor(NSensorComponent::kDq)) - ) - { - ui32 nodeId = 1; - - TString hostName = "localhost"; - TString localAddress = "::1"; - - NDqs::TServiceNodeConfig config = { - nodeId, - localAddress, - hostName, - static_cast<ui16>(interconnectPort.Addr.GetPort()), - static_cast<ui16>(grpcPort.Addr.GetPort()), - 0, // mbus - interconnectPort.Socket.Get()->Release(), - grpcPort.Socket.Get()->Release(), - 1 - }; - - ServiceNode = MakeHolder<TServiceNode>( - config, - threads, - MetricsRegistry); - - auto lwmGroup = MetricsRegistry->GetSensors()->GetSubgroup("component", "lwm"); - auto patternCache = std::make_shared<NKikimr::NMiniKQL::TComputationPatternLRUCache>(NKikimr::NMiniKQL::TComputationPatternLRUCache::Config(200_MB, 200_MB)); - NDqs::TLocalWorkerManagerOptions lwmOptions; - lwmOptions.Factory = NTaskRunnerProxy::CreateFactory(functionRegistry, compFactory, taskTransformFactory, patternCache, false); - lwmOptions.AsyncIoFactory = std::move(asyncIoFactory); - lwmOptions.FunctionRegistry = functionRegistry; - lwmOptions.TaskRunnerInvokerFactory = new NDqs::TTaskRunnerInvokerFactory(); - lwmOptions.TaskRunnerActorFactory = NDq::NTaskRunnerActor::CreateLocalTaskRunnerActorFactory( - [factory=lwmOptions.Factory](std::shared_ptr<NKikimr::NMiniKQL::TScopedAlloc> alloc, const NDq::TDqTaskSettings& task, NDqProto::EDqStatsMode statsMode, const NDq::TLogFunc& ) - { - return factory->Get(alloc, task, statsMode); - }); - lwmOptions.Counters = NDqs::TWorkerManagerCounters(lwmGroup); - lwmOptions.DropTaskCountersOnFinish = false; - auto resman = NDqs::CreateLocalWorkerManager(lwmOptions); - - ServiceNode->AddLocalService( - MakeWorkerManagerActorID(nodeId), - TActorSetupCmd(resman, TMailboxType::Simple, 0)); - - if (withSpilling) { - auto tempDir = NDq::GetTmpSpillingRootForCurrentUser(); - MakeDirIfNotExist(tempDir); - - auto spillingActor = NDq::CreateDqLocalFileSpillingService(NDq::TFileSpillingServiceConfig{.Root = tempDir, .CleanupOnShutdown = true}, MakeIntrusive<NDq::TSpillingCounters>(lwmGroup)); - - ServiceNode->AddLocalService( - NDq::MakeDqLocalFileSpillingServiceID(nodeId), - TActorSetupCmd(spillingActor, TMailboxType::Simple, 0)); - } - for (auto& [actorId, setupCmd] : additionalLocalServices) { - ServiceNode->AddLocalService(actorId, std::move(setupCmd)); - } - - auto statsCollector = CreateStatsCollector(1, *ServiceNode->GetSetup(), MetricsRegistry->GetSensors()); - - auto actorSystem = ServiceNode->StartActorSystem(); - if (metricsPusherFactory) { - actorSystem->Register(metricsPusherFactory()); - } - - actorSystem->Register(statsCollector); - - ServiceNode->StartService(dqTaskPreprocessorFactories); - } - - ~TLocalServiceHolder() - { - ServiceNode->Stop(); - } - -private: - IMetricsRegistryPtr MetricsRegistry; - THolder<TServiceNode> ServiceNode; -}; - -class TDqGatewayLocalImpl: public std::enable_shared_from_this<TDqGatewayLocalImpl> -{ - struct TRequest { - TString SessionId; - NDqs::TPlan Plan; - TVector<TString> Columns; - THashMap<TString, TString> SecureParams; - THashMap<TString, TString> GraphParams; - TDqSettings::TPtr Settings; - IDqGateway::TDqProgressWriter ProgressWriter; - THashMap<TString, TString> ModulesMapping; - bool Discard; - NThreading::TPromise<IDqGateway::TResult> Result; - ui64 ExecutionTimeout; - }; - -public: - TDqGatewayLocalImpl(THolder<TLocalServiceHolder>&& localService, const IDqGateway::TPtr& gateway) - : LocalService(std::move(localService)) - , Gateway(gateway) - , DeterministicMode(!!GetEnv("YQL_DETERMINISTIC_MODE")) - { } - - NThreading::TFuture<void> OpenSession(const TString& sessionId, const TString& username) { - return Gateway->OpenSession(sessionId, username); - } - - NThreading::TFuture<void> CloseSession(const TString& sessionId) { - return Gateway->CloseSessionAsync(sessionId); - } - - NThreading::TFuture<IDqGateway::TResult> - ExecutePlan(const TString& sessionId, NDqs::TPlan&& plan, const TVector<TString>& columns, - const THashMap<TString, TString>& secureParams, const THashMap<TString, TString>& graphParams, - const TDqSettings::TPtr& settings, - const IDqGateway::TDqProgressWriter& progressWriter, const THashMap<TString, TString>& modulesMapping, - bool discard, ui64 executionTimeout) - { - - NThreading::TFuture<IDqGateway::TResult> result; - { - TGuard<TMutex> lock(Mutex); - Queue.emplace_back(TRequest{sessionId, std::move(plan), columns, secureParams, graphParams, settings, progressWriter, modulesMapping, discard, NThreading::NewPromise<IDqGateway::TResult>(), executionTimeout}); - result = Queue.back().Result; - } - - TryExecuteNext(); - - return result; - } - - void Stop() { - Gateway->Stop(); - } - -private: - void TryExecuteNext() { - TGuard<TMutex> lock(Mutex); - if (!Queue.empty() && (!DeterministicMode || Inflight == 0)) { - auto request = std::move(Queue.front()); Queue.pop_front(); - Inflight++; - lock.Release(); - - auto weak = weak_from_this(); - - Gateway->ExecutePlan(request.SessionId, std::move(request.Plan), request.Columns, request.SecureParams, request.GraphParams, request.Settings, request.ProgressWriter, request.ModulesMapping, request.Discard, request.ExecutionTimeout) - .Apply([promise=request.Result, weak](const NThreading::TFuture<IDqGateway::TResult>& result) mutable { - try { - promise.SetValue(result.GetValue()); - } catch (...) { - promise.SetException(std::current_exception()); - } - - if (auto ptr = weak.lock()) { - { - TGuard<TMutex> lock(ptr->Mutex); - ptr->Inflight--; - } - - ptr->TryExecuteNext(); - } - }); - } - } - - THolder<TLocalServiceHolder> LocalService; - IDqGateway::TPtr Gateway; - const bool DeterministicMode; - TMutex Mutex; - TList<TRequest> Queue; - int Inflight = 0; -}; - -class TDqGatewayLocal : public IDqGateway { -public: - TDqGatewayLocal(THolder<TLocalServiceHolder>&& localService, const IDqGateway::TPtr& gateway) - : Impl(std::make_shared<TDqGatewayLocalImpl>(std::move(localService), gateway)) - {} - - ~TDqGatewayLocal() { - Impl->Stop(); - } - - NThreading::TFuture<void> OpenSession(const TString& sessionId, const TString& username) override { - return Impl->OpenSession(sessionId, username); - } - - NThreading::TFuture<void> CloseSessionAsync(const TString& sessionId) override { - return Impl->CloseSession(sessionId); - } - - NThreading::TFuture<TResult> - ExecutePlan(const TString& sessionId, NDqs::TPlan&& plan, const TVector<TString>& columns, - const THashMap<TString, TString>& secureParams, const THashMap<TString, TString>& graphParams, - const TDqSettings::TPtr& settings, - const TDqProgressWriter& progressWriter, const THashMap<TString, TString>& modulesMapping, - bool discard, ui64 executionTimeout) override - { - return Impl->ExecutePlan(sessionId, std::move(plan), columns, secureParams, graphParams, - settings, progressWriter, modulesMapping, discard, executionTimeout); - } - - void Stop() override { - Impl->Stop(); - } - -private: - std::shared_ptr<TDqGatewayLocalImpl> Impl; -}; - -THolder<TLocalServiceHolder> CreateLocalServiceHolder(const NKikimr::NMiniKQL::IFunctionRegistry* functionRegistry, - NKikimr::NMiniKQL::TComputationNodeFactory compFactory, - TTaskTransformFactory taskTransformFactory, const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories, - NBus::TBindResult interconnectPort, NBus::TBindResult grpcPort, - NDq::IDqAsyncIoFactory::TPtr asyncIoFactory, int threads, - IMetricsRegistryPtr metricsRegistry, - const std::function<IActor*(void)>& metricsPusherFactory, bool withSpilling, - TVector<std::pair<TActorId, TActorSetupCmd>>&& additionalLocalServices) -{ - return MakeHolder<TLocalServiceHolder>(functionRegistry, - compFactory, - taskTransformFactory, - dqTaskPreprocessorFactories, - interconnectPort, - grpcPort, - std::move(asyncIoFactory), - threads, - metricsRegistry, - metricsPusherFactory, - withSpilling, - std::move(additionalLocalServices)); -} - -TIntrusivePtr<IDqGateway> CreateLocalDqGateway(const NKikimr::NMiniKQL::IFunctionRegistry* functionRegistry, - NKikimr::NMiniKQL::TComputationNodeFactory compFactory, - TTaskTransformFactory taskTransformFactory, const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories, - bool withSpilling, NDq::IDqAsyncIoFactory::TPtr asyncIoFactory, int threads, - IMetricsRegistryPtr metricsRegistry, - const std::function<IActor*(void)>& metricsPusherFactory, - TVector<std::pair<TActorId, TActorSetupCmd>>&& additionalLocalServices) -{ - int startPort = 31337; - TRangeWalker<int> portWalker(startPort, startPort+100); - auto interconnectPort = BindInRange(portWalker)[1]; - auto grpcPort = BindInRange(portWalker)[1]; - - return new TDqGatewayLocal( - CreateLocalServiceHolder( - functionRegistry, - compFactory, - taskTransformFactory, - dqTaskPreprocessorFactories, - interconnectPort, - grpcPort, - std::move(asyncIoFactory), - threads, - metricsRegistry, - metricsPusherFactory, - withSpilling, - std::move(additionalLocalServices)), - CreateDqGateway("[::1]", grpcPort.Addr.GetPort())); -} - -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h b/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h deleted file mode 100644 index 30877762117..00000000000 --- a/ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h +++ /dev/null @@ -1,24 +0,0 @@ -#pragma once - -#include <ydb/library/actors/core/actorsystem.h> -#include <ydb/library/yql/providers/dq/provider/yql_dq_gateway.h> -#include <ydb/library/yql/providers/dq/interface/yql_dq_task_preprocessor.h> -#include <ydb/library/yql/dq/actors/compute/dq_compute_actor_async_io_factory.h> -#include <yql/essentials/providers/common/metrics/metrics_registry.h> - -namespace NActors { -class IActor; -} - -namespace NYql { - -TIntrusivePtr<IDqGateway> CreateLocalDqGateway(const NKikimr::NMiniKQL::IFunctionRegistry* functionRegistry, - NKikimr::NMiniKQL::TComputationNodeFactory compFactory, - TTaskTransformFactory taskTransformFactory, const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories, - bool withSpilling, - NDq::IDqAsyncIoFactory::TPtr = nullptr, int threads = 16, - IMetricsRegistryPtr metricsRegistry = {}, - const std::function<NActors::IActor*(void)>& metricsPusherFactory = {}, - TVector<std::pair<NActors::TActorId, NActors::TActorSetupCmd>>&& additionalLocalServices = {}); - -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/provider/exec/ya.make b/ydb/library/yql/providers/dq/provider/exec/ya.make index b790e06aaa4..92b15184a45 100644 --- a/ydb/library/yql/providers/dq/provider/exec/ya.make +++ b/ydb/library/yql/providers/dq/provider/exec/ya.make @@ -6,31 +6,34 @@ SRCS( ) PEERDIR( - library/cpp/yson/node - library/cpp/svnversion library/cpp/digest/md5 + library/cpp/svnversion library/cpp/threading/future - ydb/public/lib/yson_value - ydb/public/sdk/cpp/src/client/driver - yql/essentials/core - yql/essentials/core/dq_integration + library/cpp/yson/node + ydb/library/yql/dq/expr_nodes + ydb/library/yql/dq/opt ydb/library/yql/dq/runtime ydb/library/yql/dq/tasks ydb/library/yql/dq/type_ann - yql/essentials/providers/common/gateway - yql/essentials/providers/common/metrics - yql/essentials/providers/common/schema/expr - yql/essentials/providers/common/transform ydb/library/yql/providers/dq/actors - ydb/library/yql/providers/dq/api/grpc - ydb/library/yql/providers/dq/api/protos ydb/library/yql/providers/dq/common ydb/library/yql/providers/dq/counters ydb/library/yql/providers/dq/expr_nodes ydb/library/yql/providers/dq/opt ydb/library/yql/providers/dq/planner - ydb/library/yql/providers/dq/runtime + ydb/library/yql/providers/dq/provider + yql/essentials/core + yql/essentials/core/dq_integration + yql/essentials/core/peephole_opt + yql/essentials/core/services + yql/essentials/core/type_ann + yql/essentials/minikql + yql/essentials/minikql/runtime_settings + yql/essentials/providers/common/provider + yql/essentials/providers/common/schema/expr + yql/essentials/providers/common/transform yql/essentials/providers/result/expr_nodes + yql/essentials/utils/log ) YQL_LAST_ABI_VERSION() diff --git a/ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.cpp b/ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.cpp index ee06ff61468..fe6c7cebc70 100644 --- a/ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.cpp +++ b/ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.cpp @@ -20,7 +20,6 @@ #include <yql/essentials/core/dq_integration/yql_dq_integration.h> #include <ydb/library/yql/providers/dq/planner/execution_planner.h> #include <ydb/library/yql/providers/dq/provider/yql_dq_gateway.h> -#include <ydb/library/yql/providers/dq/provider/yql_dq_control.h> #include <ydb/library/yql/dq/type_ann/dq_type_ann.h> #include <ydb/library/yql/dq/runtime/dq_tasks_runner.h> @@ -507,9 +506,9 @@ private: fileLink = State->FileStorage->PutFileStripped(path, md5); } - UploadCache_->ModulesMapping.emplace(objectId + DqStrippedSuffied, path); + UploadCache_->ModulesMapping.emplace(objectId + DqStrippedSuffied(), path); - return std::make_tuple(fileLink->GetPath(), objectId + DqStrippedSuffied); + return std::make_tuple(fileLink->GetPath(), objectId + DqStrippedSuffied()); } std::tuple<TString, TString> GetPathAndObjectId(const TFilePathWithMd5& pathWithMd5) const { diff --git a/ydb/library/yql/providers/dq/provider/ut/ya.make b/ydb/library/yql/providers/dq/provider/ut/ya.make deleted file mode 100644 index 5efa4874a0f..00000000000 --- a/ydb/library/yql/providers/dq/provider/ut/ya.make +++ /dev/null @@ -1,41 +0,0 @@ -UNITTEST_FOR(ydb/library/yql/providers/dq/provider) - -SRCS( - yql_dq_provider_ut.cpp -) - -PEERDIR( - ydb/library/yql/dq/actors/compute - ydb/library/yql/dq/comp_nodes - ydb/library/yql/dq/transform - ydb/library/yql/providers/dq/local_gateway - ydb/library/yql/providers/dq/provider - ydb/library/yql/providers/dq/provider/exec - library/cpp/lwtrace - library/cpp/lwtrace/mon - library/cpp/testing/unittest - yt/yql/providers/yt/codec/codegen - yt/yql/providers/yt/comp_nodes/llvm16 - yt/yql/providers/yt/gateway/file - yt/yql/providers/yt/lib/ut_common - yt/yql/providers/yt/provider - yql/essentials/core/cbo/simple - yql/essentials/core/facade - yql/essentials/core/file_storage - yql/essentials/core/services/mounts - yql/essentials/minikql/comp_nodes/llvm16 - yql/essentials/providers/common/comp_nodes - yql/essentials/public/udf/service/exception_policy - yql/essentials/sql/pg -) - -YQL_LAST_ABI_VERSION() - -IF (SANITIZER_TYPE) - SIZE(LARGE) - INCLUDE(${ARCADIA_ROOT}/ydb/tests/large.inc) -ELSE() - SIZE(MEDIUM) -ENDIF() - -END() diff --git a/ydb/library/yql/providers/dq/provider/ya.make b/ydb/library/yql/providers/dq/provider/ya.make index 5a1d26fbbc1..7364d9039bc 100644 --- a/ydb/library/yql/providers/dq/provider/ya.make +++ b/ydb/library/yql/providers/dq/provider/ya.make @@ -1,15 +1,12 @@ LIBRARY() SRCS( - yql_dq_control.cpp - yql_dq_control.h yql_dq_datasink_constraints.cpp yql_dq_datasink_type_ann.cpp yql_dq_datasink_type_ann.h yql_dq_datasource_constraints.cpp yql_dq_datasource_type_ann.cpp yql_dq_datasource_type_ann.h - yql_dq_gateway.cpp yql_dq_gateway.h yql_dq_provider.cpp yql_dq_provider.h @@ -27,44 +24,40 @@ SRCS( ) PEERDIR( - library/cpp/threading/task_scheduler - library/cpp/threading/future - library/cpp/svnversion - library/cpp/yson/node library/cpp/yson ydb/library/yql/dq/constraints ydb/library/yql/dq/expr_nodes - ydb/library/yql/dq/tasks + ydb/library/yql/dq/opt ydb/library/yql/dq/transform ydb/library/yql/dq/type_ann - ydb/library/yql/providers/dq/actors - ydb/library/yql/providers/dq/api/grpc + ydb/library/yql/providers/common/http_gateway ydb/library/yql/providers/dq/api/protos ydb/library/yql/providers/dq/common - ydb/library/yql/providers/dq/config ydb/library/yql/providers/dq/expr_nodes ydb/library/yql/providers/dq/opt ydb/library/yql/providers/dq/planner - ydb/public/lib/yson_value ydb/public/sdk/cpp/src/client/driver - ydb/public/sdk/cpp/src/library/grpc/client - yql/essentials/providers/result/expr_nodes yql/essentials/ast yql/essentials/core + yql/essentials/core/cbo yql/essentials/core/dq_integration yql/essentials/core/dq_integration/transform + yql/essentials/core/expr_nodes + yql/essentials/core/file_storage yql/essentials/core/issue + yql/essentials/core/services + yql/essentials/core/type_ann + yql/essentials/minikql + yql/essentials/minikql/computation yql/essentials/providers/common/activation yql/essentials/providers/common/config/transformer yql/essentials/providers/common/gateway yql/essentials/providers/common/metrics yql/essentials/providers/common/proto - yql/essentials/providers/common/schema/expr + yql/essentials/providers/common/provider yql/essentials/providers/common/transform - yql/essentials/minikql - yql/essentials/public/issue - yql/essentials/utils/backtrace - yql/essentials/utils/failure_injector + yql/essentials/providers/result/expr_nodes + yql/essentials/utils/log ) YQL_LAST_ABI_VERSION() @@ -74,7 +67,3 @@ END() RECURSE( exec ) - -RECURSE_FOR_TESTS( - ut -) diff --git a/ydb/library/yql/providers/dq/provider/yql_dq_control.cpp b/ydb/library/yql/providers/dq/provider/yql_dq_control.cpp deleted file mode 100644 index 847fcfb32fc..00000000000 --- a/ydb/library/yql/providers/dq/provider/yql_dq_control.cpp +++ /dev/null @@ -1,197 +0,0 @@ -#include "yql_dq_control.h" - -#include <ydb/library/yql/providers/dq/api/grpc/api.grpc.pb.h> -#include <ydb/library/yql/providers/dq/config/config.pb.h> -#include <yql/essentials/utils/log/log.h> -#include <yql/essentials/minikql/mkql_function_registry.h> - -#include <ydb/public/lib/yson_value/ydb_yson_value.h> - -#include <ydb/public/sdk/cpp/src/library/grpc/client/grpc_client_low.h> - -#include <library/cpp/svnversion/svnversion.h> - -#include <util/memory/blob.h> -#include <util/string/builder.h> -#include <util/system/file.h> - -namespace NYql { - -using TFileResource = Yql::DqsProto::TFile; - -const TString DqStrippedSuffied = ".s"; - -class TDqControl : public IDqControl { - -public: - TDqControl(const NYdbGrpc::TGRpcClientConfig &grpcConf, int threads, const TVector<TFileResource> &files) - : GrpcClient(threads) - , Service(GrpcClient.CreateGRpcServiceConnection<Yql::DqsProto::DqService>(grpcConf)) - , Files(files) - { } - - // after call, forking process in not allowed - bool IsReady(const TMap<TString, TString>& additinalFiles) override { - Yql::DqsProto::IsReadyRequest request; - for (const auto& file : Files) { - *request.AddFiles() = file; - } - - for (const auto& [path, objectId] : additinalFiles){ - TFileResource r; - r.SetLocalPath(path); - r.SetObjectType(Yql::DqsProto::TFile::EUDF_FILE); - r.SetObjectId(objectId); - r.SetSize(TFile(path, OpenExisting | RdOnly).GetLength()); - *request.AddFiles() = r; - } - - auto promise = NThreading::NewPromise<bool>(); - auto callback = [promise](NYdbGrpc::TGrpcStatus&& status, Yql::DqsProto::IsReadyResponse&& resp) mutable { - Y_UNUSED(resp); - - promise.SetValue(status.Ok() && resp.GetIsReady()); - }; - - NYdbGrpc::TCallMeta meta; - meta.Timeout = std::chrono::seconds(1); - Service->DoRequest<Yql::DqsProto::IsReadyRequest, Yql::DqsProto::IsReadyResponse>( - request, callback, &Yql::DqsProto::DqService::Stub::AsyncIsReady, meta); - - try { - return promise.GetFuture().GetValueSync(); - } catch (...) { - YQL_CLOG(INFO, ProviderDq) << "DqControl IsReady Exception " << CurrentExceptionMessage(); - return false; - } - } - -private: - NYdbGrpc::TGRpcClientLow GrpcClient; - std::unique_ptr<NYdbGrpc::TServiceConnection<Yql::DqsProto::DqService>> Service; - const TVector<TFileResource> &Files; -}; - - -class TDqControlFactory : public IDqControlFactory { - -public: - TDqControlFactory( - const TString &host, - int port, - int threads, - const TMap<TString, TString>& udfs, - const TString& vanillaLitePath, - const TString& vanillaLiteMd5, - const THashSet<TString> &filter, - bool enableStrip, - const TFileStoragePtr& fileStorage - ) - : Threads(threads) - , GrpcConf(TStringBuilder() << host << ":" << port) - , IndexedUdfFilter(filter) - , EnableStrip(enableStrip) - , FileStorage(fileStorage) - { - if (!vanillaLitePath.empty()) { - TString path = vanillaLitePath; - TString objectId = GetProgramCommitId(); - - TString newPath, newObjectId; - std::tie(newPath, newObjectId) = GetPathAndObjectId(path, objectId, vanillaLiteMd5); - - TFileResource vanillaLite; - vanillaLite.SetLocalPath(newPath); - vanillaLite.SetName(vanillaLitePath.substr(vanillaLitePath.rfind('/') + 1)); - vanillaLite.SetObjectType(Yql::DqsProto::TFile::EEXE_FILE); - vanillaLite.SetObjectId(newObjectId); - vanillaLite.SetSize(TFile(newPath, OpenExisting | RdOnly).GetLength()); - Files.push_back(vanillaLite); - } - - for (const auto& [path, objectId] : udfs){ - YQL_CLOG(DEBUG, ProviderDq) << "DQ control, adding file: " << path << " with objectId " << objectId; - TString newPath, newObjectId; - std::tie(newPath, newObjectId) = GetPathAndObjectId(path, objectId, objectId); - - YQL_CLOG(DEBUG, ProviderDq) << "DQ control, rewrite path/objectId: " << newPath << ", " << newObjectId; - TFileResource r; - r.SetLocalPath(newPath); - r.SetObjectType(Yql::DqsProto::TFile::EUDF_FILE); - r.SetObjectId(newObjectId); - r.SetSize(TFile(newPath, OpenExisting | RdOnly).GetLength()); - Files.push_back(r); - } - } - - IDqControlPtr GetControl() override { - return new TDqControl(GrpcConf, Threads, Files); - } - - const THashSet<TString>& GetIndexedUdfFilter() override { - return IndexedUdfFilter; - } - - bool StripEnabled() const override { - return EnableStrip; - } - -private: - std::tuple<TString, TString> GetPathAndObjectId(const TString& path, const TString& objectId, const TString& md5 = {}) { - if (!EnableStrip) { - return std::make_tuple(path, objectId); - } - - TFileLinkPtr& fileLink = FileLinks[objectId]; - if (!fileLink) { - fileLink = FileStorage->PutFileStripped(path, md5); - } - - return std::make_tuple(fileLink->GetPath(), objectId + DqStrippedSuffied); - } - - int Threads; - TVector<TFileResource> Files; - NYdbGrpc::TGRpcClientConfig GrpcConf; - THashSet<TString> IndexedUdfFilter; - THashMap<TString, TFileLinkPtr> FileLinks; - bool EnableStrip; - const TFileStoragePtr FileStorage; -}; - -IDqControlFactoryPtr CreateDqControlFactory(const NProto::TDqConfig& config, const TMap<TString, TString>& udfs, const TFileStoragePtr& fileStorage) { - THashSet<TString> indexedUdfFilter(config.GetControl().GetIndexedUdfsToWarmup().begin(), config.GetControl().GetIndexedUdfsToWarmup().end()); - return CreateDqControlFactory( - config.GetPort(), - config.GetYtBackends()[0].GetVanillaJobLite(), - config.GetYtBackends()[0].GetVanillaJobLiteMd5(), - config.GetControl().GetEnableStrip(), - indexedUdfFilter, - udfs, - fileStorage - ); -} - -IDqControlFactoryPtr CreateDqControlFactory( - const uint32_t port, - const TString& vanillaJobLite, - const TString& vanillaJobLiteMd5, - const bool enableStrip, - const THashSet<TString> indexedUdfFilter, - const TMap<TString, TString>& udfs, - const TFileStoragePtr& fileStorage) -{ - return new TDqControlFactory( - "localhost", - port, - 2, - udfs, - vanillaJobLite, - vanillaJobLiteMd5, - indexedUdfFilter, - enableStrip, - fileStorage - ); -} - -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/provider/yql_dq_control.h b/ydb/library/yql/providers/dq/provider/yql_dq_control.h deleted file mode 100644 index 6335449199c..00000000000 --- a/ydb/library/yql/providers/dq/provider/yql_dq_control.h +++ /dev/null @@ -1,46 +0,0 @@ -#pragma once - -#include <util/generic/ptr.h> -#include <util/generic/map.h> - -#include <yql/essentials/core/file_storage/file_storage.h> - -namespace NYql { - -namespace NProto { -class TDqConfig; -} - -class IDqControl : public TThrRefBase { -public: - using TPtr = TIntrusivePtr<IDqControl>; - - virtual bool IsReady(const TMap<TString, TString>& udfs = TMap<TString, TString>())= 0; -}; - -using IDqControlPtr = TIntrusivePtr<IDqControl>; - -class IDqControlFactory : public TThrRefBase { -public: - using TPtr = TIntrusivePtr<IDqControlFactory>; - - virtual IDqControlPtr GetControl() = 0; - virtual const THashSet<TString>& GetIndexedUdfFilter() = 0; - virtual bool StripEnabled() const = 0; -}; - -using IDqControlFactoryPtr = TIntrusivePtr<IDqControlFactory>; - -IDqControlFactoryPtr CreateDqControlFactory(const NProto::TDqConfig& config, const TMap<TString, TString>& udfs, const TFileStoragePtr& fileStorage); -IDqControlFactoryPtr CreateDqControlFactory( - const uint32_t port, - const TString& vanillaJobLite, - const TString& vanillaJobLiteMd5, - const bool enableStrip, - const THashSet<TString> indexedUdfFilter, - const TMap<TString, TString>& udfs, - const TFileStoragePtr& fileStorage); - -extern const TString DqStrippedSuffied; - -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/provider/yql_dq_gateway.cpp b/ydb/library/yql/providers/dq/provider/yql_dq_gateway.cpp deleted file mode 100644 index fadd62270f1..00000000000 --- a/ydb/library/yql/providers/dq/provider/yql_dq_gateway.cpp +++ /dev/null @@ -1,802 +0,0 @@ -#include "yql_dq_gateway.h" - -#include <yql/essentials/providers/common/provider/yql_provider_names.h> -#include <ydb/library/yql/providers/dq/api/grpc/api.grpc.pb.h> -#include <ydb/library/yql/providers/dq/common/yql_dq_common.h> -#include <ydb/library/yql/providers/dq/actors/proto_builder.h> -#include <yql/essentials/utils/backtrace/backtrace.h> -#include <yql/essentials/utils/failure_injector/failure_injector.h> -#include <yql/essentials/public/issue/yql_issue_message.h> -#include <ydb/library/yql/providers/dq/config/config.pb.h> -#include <yql/essentials/utils/log/log.h> - -#include <ydb/public/lib/yson_value/ydb_yson_value.h> - -#include <ydb/public/sdk/cpp/src/library/grpc/client/grpc_client_low.h> - -#include <library/cpp/yson/node/node_io.h> -#include <library/cpp/threading/task_scheduler/task_scheduler.h> - -#include <util/system/mutex.h> -#include <util/generic/hash.h> -#include <util/string/builder.h> - -#include <utility> - -namespace NYql { - -using namespace NThreading; - -class TPlanPrinter { -public: - TStringBuilder b; - - void DescribeChannel(const auto& ch, bool spilling) { - if (spilling) { - b << "Ch" << ch.GetId() << " [shape=diamond, label=\"Ch" << ch.GetId() << "\", color=\"red\"];"; - } else { - b << "Ch" << ch.GetId() << " [shape=diamond, label=\"Ch" << ch.GetId() << "\"];"; - } - } - - void PrintInputChannel(const auto& ch, const auto& type) { - b << "Ch" << ch.GetId() << " -> T" << ch.GetDstTaskId() << " [label=" << "\"" << type << "\"];\n"; - } - - void PrintOutputChannel(const auto& ch, const auto& type) { - b << "T" << ch.GetSrcTaskId() << " -> Ch" << ch.GetId() << " [label=" << "\"" << type << "\"];\n"; - } - - void PrintSource(auto taskId, auto sourceIndex) { - b << "S" << taskId << "_" << sourceIndex << " -> T" << taskId << " [label=" << "\"S" << sourceIndex << "\"];\n"; - } - - void DescribeSource(auto taskId, auto sourceIndex) { - b << "S" << taskId << "_" << sourceIndex << " "; - b << "[shape=square, label=\"" << taskId << "/" << sourceIndex << "\"];\n"; - } - - void PrintTask(const auto& task) { - int index = 0; - for (const auto& input : task.GetInputs()) { - TString inputName = "Unknown"; - bool isSource = false; - if (input.HasUnionAll()) { inputName = "UnionAll"; } - else if (input.HasMerge()) { inputName = "Merge"; } - else if (input.HasSource()) { inputName = "Source"; isSource = true; } - if (isSource) { - PrintSource(task.GetId(), index); - } else { - for (const auto& ch : input.GetChannels()) { - PrintInputChannel(ch, inputName); - } - } - index ++; - } - for (const auto& output : task.GetOutputs()) { - TString outputName = "Unknown"; - if (output.HasMap()) { outputName = "Map"; } - else if (output.HasRangePartition()) { outputName = "Range"; } - else if (output.HasHashPartition()) { outputName = "Hash"; } - else if (output.HasBroadcast()) { outputName = "Broadcast"; } - // TODO: effects, sink - for (const auto& ch : output.GetChannels()) { - PrintOutputChannel(ch, outputName); - } - } - } - - void DescribeTask(const auto& task) { - b << "T" << task.GetId() << " [shape=circle, label=\"" << task.GetId() << "/" << task.GetStageId() << "\"];\n"; - int index = 0; - for (const auto& input : task.GetInputs()) { - if (input.HasSource()) { - DescribeSource(task.GetId(), index); - } - index ++; - } - for (const auto& output : task.GetOutputs()) { - for (const auto& ch : output.GetChannels()) { - DescribeChannel(ch, task.GetEnableSpilling()); - } - } - } - - TString Print(const NDqs::TPlan& plan) { - b.clear(); - b << "digraph G {\n"; - for (const auto& task : plan.Tasks) { - DescribeTask(task); - } - b << "\n"; - for (const auto& task : plan.Tasks) { - PrintTask(task); - } - b << "}\n"; - return b; - } -}; - -class TDqTaskScheduler : public TTaskScheduler { -private: - struct TDelay: public TTaskScheduler::ITask { - TDelay(TPromise<void> p) - : Promise(std::move(p)) - { } - - TInstant Process() override { - Promise.SetValue(); - return TInstant::Max(); - } - - TPromise<void> Promise; - }; - -public: - TDqTaskScheduler() - : TTaskScheduler(1) // threads - {} - - TFuture<void> Delay(TDuration duration) { - TPromise<void> promise = NewPromise(); - - auto future = promise.GetFuture(); - - if (!Add(MakeIntrusive<TDelay>(promise), TInstant::Now() + duration)) { - promise.SetException("cannot delay"); - } - - return future; - } -}; - -class TDqGatewaySession: public std::enable_shared_from_this<TDqGatewaySession> { -public: - using TResult = IDqGateway::TResult; - using TDqProgressWriter = IDqGateway::TDqProgressWriter; - - TDqGatewaySession(const TString& sessionId, TDqTaskScheduler& taskScheduler, NYdbGrpc::TServiceConnection<Yql::DqsProto::DqService>& service, TFuture<void>&& openSessionFuture) - : SessionId(sessionId) - , TaskScheduler(taskScheduler) - , Service(service) - , OpenSessionFuture(std::move(openSessionFuture)) - { - } - - const TString& GetSessionId() const { - return SessionId; - } - - template<typename RespType> - void OnResponse(TPromise<TResult> promise, NYdbGrpc::TGrpcStatus&& status, RespType&& resp, const NCommon::TResultFormatSettings& resultFormatSettings, const THashMap<TString, TString>& modulesMapping, bool alwaysFallback = false) { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(SessionId); - YQL_CLOG(TRACE, ProviderDq) << "TDqGateway::callback"; - - TResult result; - - bool error = false; - bool fallback = false; - result.Timeout = resp.GetTimeout(); - - if (status.Ok()) { - YQL_CLOG(TRACE, ProviderDq) << "TDqGateway::Ok"; - - result.Truncated = resp.GetTruncated(); - - TOperationStatistics statistics; - - for (const auto& t : resp.GetMetric()) { - YQL_CLOG(TRACE, ProviderDq) << "Counter: " << t.GetName() << " : " << t.GetSum() << " : " << t.GetCount(); - TOperationStatistics::TEntry entry( - t.GetName(), - t.GetSum(), - t.GetMax(), - t.GetMin(), - t.GetAvg(), - t.GetCount()); - statistics.Entries.push_back(entry); - } - - result.Statistics = statistics; - - NYql::TIssues issues; - auto operation = resp.operation(); - - for (auto& message_ : *operation.Mutableissues()) { - TDeque<std::remove_reference_t<decltype(message_)>*> queue; - queue.push_front(&message_); - while (!queue.empty()) { - auto& message = *queue.front(); - queue.pop_front(); - message.Setmessage(NBacktrace::Symbolize(message.Getmessage(), modulesMapping)); - for (auto& subMsg : *message.Mutableissues()) { - queue.push_back(&subMsg); - } - } - } - - NYql::IssuesFromMessage(operation.issues(), issues); - error = false; - for (const auto& issue : issues) { - if (issue.GetSeverity() <= TSeverityIds::S_ERROR) { - error = true; - } - if (issue.GetCode() == TIssuesIds::DQ_GATEWAY_NEED_FALLBACK_ERROR) { - fallback = true; - } - } - - // TODO: Save statistics in case of result failure - if (!error) { - Yql::DqsProto::ExecuteQueryResult queryResult; - resp.operation().result().UnpackTo(&queryResult); - TVector<NDq::TDqSerializedBatch> rows; - for (const auto& s : queryResult.Getsample()) { - NDq::TDqSerializedBatch batch; - batch.Proto = s; - rows.emplace_back(std::move(batch)); - } - - result.AddIssues(issues); - try { - NYql::NDqs::TProtoBuilder protoBuilder(resultFormatSettings.ResultType, resultFormatSettings.Columns); - - bool ysonTruncated = false; - result.Data = protoBuilder.BuildYson(std::move(rows), - result.Truncated ? resultFormatSettings.SizeLimit.GetOrElse(Max<ui64>()) : Max<ui64>(), - result.Truncated ? resultFormatSettings.RowsLimit.GetOrElse(Max<ui64>()) : Max<ui64>(), - &ysonTruncated); - - result.Truncated = result.Truncated || ysonTruncated; - result.SetSuccess(); - } catch (...) { - YQL_CLOG(ERROR, ProviderDq) << "Failed to build yson result: " << CurrentExceptionMessage(); - error = true; - auto issue = TIssue("Failed to build query result (probably due to malformed UDF)"); - result.AddIssue(issue.SetCode(TIssuesIds::DQ_GATEWAY_ERROR, TSeverityIds::S_ERROR)); - } - } else { - YQL_CLOG(ERROR, ProviderDq) << "Issue " << issues.ToString(); - result.AddIssues(issues); - if (fallback) { - result.Fallback = true; - result.SetSuccess(); - } - } - } else { - YQL_CLOG(ERROR, ProviderDq) << "Issue " << status.Msg; - auto issue = TIssue(TStringBuilder{} << "Error " << status.GRpcStatusCode << " message: " << status.Msg); - result.Retriable = status.GRpcStatusCode == grpc::CANCELLED; - if ((status.GRpcStatusCode == grpc::UNAVAILABLE /* terminating state */ - || status.GRpcStatusCode == grpc::CANCELLED /* server crashed or stopped before task process */) - || status.GRpcStatusCode == grpc::RESOURCE_EXHAUSTED /* send message limit */ - || status.GRpcStatusCode == grpc::INVALID_ARGUMENT /* Bad session */ - ) - { - YQL_CLOG(ERROR, ProviderDq) << "Fallback " << status.GRpcStatusCode; - result.Fallback = true; - result.SetSuccess(); - result.AddIssue(issue.SetCode(TIssuesIds::DQ_GATEWAY_NEED_FALLBACK_ERROR, TSeverityIds::S_ERROR)); - } else { - error = true; - result.AddIssue(issue.SetCode(TIssuesIds::DQ_GATEWAY_ERROR, TSeverityIds::S_ERROR)); - } - } - - if (error && alwaysFallback) { - YQL_CLOG(ERROR, ProviderDq) << "Force Fallback"; - result.Fallback = true; - result.ForceFallback = true; - result.SetSuccess(); - } - - promise.SetValue(result); - } - - template <typename TResponse, typename TRequest, typename TStub> - TFuture<TResult> WithRetry( - const TRequest& queryPB, - TStub stub, - int retry, - const TDqSettings::TPtr& settings, - const NCommon::TResultFormatSettings& resultFormatSettings, - const THashMap<TString, TString>& modulesMapping, - const TDqProgressWriter& progressWriter - ) { - auto backoff = TDuration::MilliSeconds(settings->RetryBackoffMs.Get().GetOrElse(1000)); - auto promise = NewPromise<TResult>(); - const auto fallbackPolicy = settings->FallbackPolicy.Get().GetOrElse(EFallbackPolicy::Default); - const auto alwaysFallback = EFallbackPolicy::Always == fallbackPolicy; - auto self = weak_from_this(); - auto callback = [self, promise, sessionId = SessionId, alwaysFallback, resultFormatSettings, modulesMapping]( - NYdbGrpc::TGrpcStatus&& status, TResponse&& resp) mutable { - auto this_ = self.lock(); - if (!this_) { - YQL_CLOG(DEBUG, ProviderDq) << "Session was closed: " << sessionId; - promise.SetException("Session was closed"); - return; - } - - this_->OnResponse(std::move(promise), std::move(status), std::move(resp), resultFormatSettings, - modulesMapping, alwaysFallback); - }; - - Service.DoRequest<TRequest, TResponse>(queryPB, callback, stub); - - ScheduleQueryStatusRequest(progressWriter, queryPB.GetQuerySeqNo()); - - return promise.GetFuture().Apply([=](const TFuture<TResult>& result) { - if (result.HasException()) { - return result; - } - auto value = result.GetValue(); - auto this_ = self.lock(); - - if (value.Success() || retry == 0 || !value.Retriable || !this_) { - return result; - } - - return this_->TaskScheduler.Delay(backoff) - .Apply([=, sessionId = this_->GetSessionId()](const TFuture<void>& result) { - auto this_ = self.lock(); - try { - result.TryRethrow(); - if (!this_) { - YQL_CLOG(DEBUG, ProviderDq) << "Session was closed: " << sessionId; - throw std::runtime_error("Session was closed"); - } - } catch (...) { - return MakeErrorFuture<TResult>(std::current_exception()); - } - return this_->WithRetry<TResponse>(queryPB, stub, retry - 1, settings, resultFormatSettings, - modulesMapping, progressWriter); - }); - }); - } - - TFuture<TResult> - ExecutePlan(NDqs::TPlan&& plan, const TVector<TString>& columns, - const THashMap<TString, TString>& secureParams, const THashMap<TString, TString>& graphParams, - const TDqSettings::TPtr& settings, - const TDqProgressWriter& progressWriter, const THashMap<TString, TString>& modulesMapping, - bool discard, ui64 executionTimeout) - { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(SessionId); - - Yql::DqsProto::ExecuteGraphRequest queryPB; - for (const auto& task : plan.Tasks) { - auto* t = queryPB.AddTask(); - *t = task; - - Yql::DqsProto::TTaskMeta taskMeta; - task.GetMeta().UnpackTo(&taskMeta); - - for (auto& file : taskMeta.GetFiles()) { - YQL_ENSURE(!file.GetObjectId().empty()); - } - } - queryPB.SetExecutionTimeout(executionTimeout); - queryPB.SetSession(SessionId); - queryPB.SetResultType(plan.ResultType); - queryPB.SetSourceId(plan.SourceID.NodeId()-1); - for (const auto& column : columns) { - *queryPB.AddColumns() = column; - } - settings->Save(queryPB); - - NCommon::TResultFormatSettings resultFormatSettings; - resultFormatSettings.Columns = columns; - resultFormatSettings.ResultType = plan.ResultType; - resultFormatSettings.SizeLimit = settings->_AllResultsBytesLimit.Get(); - resultFormatSettings.RowsLimit = settings->_RowsLimitPerWrite.Get(); - - YQL_CLOG(TRACE, ProviderDq) << TPlanPrinter().Print(plan); - - { - auto& secParams = *queryPB.MutableSecureParams(); - for (const auto&[k, v] : secureParams) { - secParams[k] = v; - } - } - - { - auto& gParams = *queryPB.MutableGraphParams(); - for (const auto&[k, v] : graphParams) { - gParams[k] = v; - } - } - - queryPB.SetDiscard(discard); - queryPB.SetQuerySeqNo(QuerySeqNo++); - - int retry = settings->MaxRetries.Get().GetOrElse(5); - - YQL_CLOG(DEBUG, ProviderDq) << "Send query of size " << queryPB.ByteSizeLong(); - - auto self = weak_from_this(); - return OpenSessionFuture.Apply([self, sessionId = SessionId, queryPB, retry, settings, resultFormatSettings, modulesMapping, - progressWriter](const TFuture<void>& f) { - f.TryRethrow(); - auto this_ = self.lock(); - if (!this_) { - YQL_CLOG(DEBUG, ProviderDq) << "Session was closed: " << sessionId; - return MakeErrorFuture<TResult>(std::make_exception_ptr(std::runtime_error("Session was closed"))); - } - - return this_->WithRetry<Yql::DqsProto::ExecuteGraphResponse>( - queryPB, - &Yql::DqsProto::DqService::Stub::AsyncExecuteGraph, - retry, - settings, - resultFormatSettings, - modulesMapping, - progressWriter); - }); - } - - TFuture<void> Close() { - Yql::DqsProto::CloseSessionRequest request; - request.SetSession(SessionId); - - auto promise = NewPromise<void>(); - auto callback = [promise, sessionId = SessionId](NYdbGrpc::TGrpcStatus&& status, Yql::DqsProto::CloseSessionResponse&& resp) mutable { - Y_UNUSED(resp); - YQL_LOG_CTX_ROOT_SESSION_SCOPE(sessionId); - if (status.Ok()) { - YQL_CLOG(DEBUG, ProviderDq) << "Async close session OK"; - promise.SetValue(); - } else { - YQL_CLOG(ERROR, ProviderDq) << "Async close session error: " << status.GRpcStatusCode << ", message: " << status.Msg; - promise.SetException(TStringBuilder() << "Async close session error: " << status.GRpcStatusCode << ", message: " << status.Msg); - } - }; - - Service.DoRequest<Yql::DqsProto::CloseSessionRequest, Yql::DqsProto::CloseSessionResponse>( - request, callback, &Yql::DqsProto::DqService::Stub::AsyncCloseSession); - return promise.GetFuture(); - } - - void OnRequestQueryStatus(const TDqProgressWriter& progressWriter, IDqGateway::TProgressWriterState state, bool ok, uint64_t querySeqNo) { - if (ok) { - ScheduleQueryStatusRequest(progressWriter, querySeqNo); - if (!state.empty()) { - progressWriter(std::move(state)); - } - } - } - - static std::unordered_map<ui64, IDqGateway::TStageStats> ExtractStats(const Yql::DqsProto::QueryStatusResponse& resp) { - std::unordered_map<ui64, IDqGateway::TStageStats> ret; - for (const auto& metric : resp.GetMetric()) { - auto longName = metric.GetName(); - TString prefix; - TString name; - std::map<TString, TString> labels; - if (!NYql::NCommon::ParseCounterName(&prefix, &labels, &name, longName)) { - continue; - } - - auto maybeStage = labels.find("Stage"); - if (maybeStage == labels.end()) { - continue; - } - auto stageId = atoi(maybeStage->second.data()); - if (!stageId) { - continue; - } - auto& stage = ret[stageId]; - - if (name == "OutputRows") { - stage.OutputRows += metric.GetSum(); - } - if (name == "InputRows") { - stage.InputRows += metric.GetSum(); - } - if (name == "OutputBytes") { - stage.OutputBytes += metric.GetSum(); - } - if (name == "InputBytes") { - stage.InputBytes += metric.GetSum(); - } - } - return ret; - } - - void RequestQueryStatus(const TDqProgressWriter& progressWriter, uint64_t querySeqNo) { - Yql::DqsProto::QueryStatusRequest request; - request.SetSession(SessionId); - request.SetQuerySeqNo(querySeqNo); - auto self = weak_from_this(); - auto callback = [self, progressWriter, querySeqNo](NYdbGrpc::TGrpcStatus&& status, Yql::DqsProto::QueryStatusResponse&& resp) { - auto this_ = self.lock(); - if (!this_) { - return; - } - - this_->OnRequestQueryStatus(progressWriter, std::move(IDqGateway::TProgressWriterState{resp.GetStatus(), std::move(ExtractStats(resp))}), status.Ok(), querySeqNo); - }; - - Service.DoRequest<Yql::DqsProto::QueryStatusRequest, Yql::DqsProto::QueryStatusResponse>( - request, callback, &Yql::DqsProto::DqService::Stub::AsyncQueryStatus, {}, nullptr); - } - - void ScheduleQueryStatusRequest(const TDqProgressWriter& progressWriter, uint64_t querySeqNo) { - auto self = weak_from_this(); - TaskScheduler.Delay(TDuration::MilliSeconds(1000)).Subscribe([self, progressWriter, querySeqNo](const TFuture<void>& f) { - auto this_ = self.lock(); - if (!this_) { - return; - } - - if (!f.HasException()) { - this_->RequestQueryStatus(progressWriter, querySeqNo); - } - }); - } - -private: - const TString SessionId; - TDqTaskScheduler& TaskScheduler; - NYdbGrpc::TServiceConnection<Yql::DqsProto::DqService>& Service; - - TMutex ProgressMutex; - - std::optional<TDqProgressWriter> ProgressWriter; - TString Status; - TFuture<void> OpenSessionFuture; - std::atomic<ui64> QuerySeqNo = 1; -}; - -class TDqGatewayImpl: public std::enable_shared_from_this<TDqGatewayImpl> { - using TResult = IDqGateway::TResult; - using TDqProgressWriter = IDqGateway::TDqProgressWriter; - -public: - TDqGatewayImpl(const TString& host, int port, TDuration timeout = TDuration::Minutes(60), TDuration requestTimeout = TDuration::Max()) - : GrpcConf(TStringBuilder() << host << ":" << port, requestTimeout) - , GrpcClient(1) - , Service(GrpcClient.CreateGRpcServiceConnection<Yql::DqsProto::DqService>(GrpcConf)) - , TaskScheduler() - , OpenSessionTimeout(timeout) - , IsStopped(false) - { - TaskScheduler.Start(); - } - - ~TDqGatewayImpl() { - Stop(); - } - - void Stop() { - bool expected = false; - if (!IsStopped.compare_exchange_strong(expected, true)) { - return; - } - - decltype(Sessions) sessions; - with_lock (Mutex) { - sessions = std::move(Sessions); - } - for (auto& pair: sessions) { - try { - pair.second->Close().GetValueSync(); - } catch (...) { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(pair.first); - YQL_CLOG(ERROR, ProviderDq) << "Error closing session " << pair.first << ": " << CurrentExceptionMessage(); - } - } - sessions.clear(); // Destroy session objects explicitly before stopping grpc - TaskScheduler.Stop(); - try { - GrpcClient.Stop(/* wait = */ true); - } catch (...) { - YQL_CLOG(ERROR, ProviderDq) << "Error while stopping GRPC client: " << CurrentExceptionMessage(); - } - } - - void DropSession(const TString& sessionId) { - with_lock (Mutex) { - Sessions.erase(sessionId); - } - } - - TFuture<void> OpenSession(const TString& sessionId, const TString& username) { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(sessionId); - YQL_CLOG(INFO, ProviderDq) << "OpenSession"; - - auto promise = NewPromise<void>(); - std::shared_ptr<TDqGatewaySession> session = std::make_shared<TDqGatewaySession>(sessionId, TaskScheduler, *Service, promise.GetFuture()); - with_lock (Mutex) { - if (!Sessions.emplace(sessionId, session).second) { - return MakeErrorFuture<void>(std::make_exception_ptr(yexception() << "Duplicate session id: " << sessionId)); - } - } - - Yql::DqsProto::OpenSessionRequest request; - request.SetSession(sessionId); - request.SetUsername(username); - - NYdbGrpc::TCallMeta meta; - meta.Timeout = OpenSessionTimeout ? NYdb::TDeadline::SafeDurationCast(OpenSessionTimeout) : NYdb::TDeadline::Duration::max(); - - auto self = weak_from_this(); - auto callback = [self, promise, sessionId](NYdbGrpc::TGrpcStatus&& status, Yql::DqsProto::OpenSessionResponse&& resp) mutable { - Y_UNUSED(resp); - YQL_LOG_CTX_ROOT_SESSION_SCOPE(sessionId); - auto this_ = self.lock(); - if (!this_) { - YQL_CLOG(ERROR, ProviderDq) << "Session was closed: " << sessionId; - promise.SetException("Session was closed"); - return; - } - if (status.Ok()) { - YQL_CLOG(INFO, ProviderDq) << "OpenSession OK"; - this_->SchedulePingSessionRequest(sessionId); - promise.SetValue(); - } else { - YQL_CLOG(ERROR, ProviderDq) << "OpenSession error: " << status.Msg; - this_->DropSession(sessionId); - promise.SetException(TString{status.Msg}); - } - }; - - Service->DoRequest<Yql::DqsProto::OpenSessionRequest, Yql::DqsProto::OpenSessionResponse>( - request, callback, &Yql::DqsProto::DqService::Stub::AsyncOpenSession, meta); - - return MakeFuture(); - } - - void SchedulePingSessionRequest(const TString& sessionId) { - auto self = weak_from_this(); - auto callback = [self, sessionId] (NYdbGrpc::TGrpcStatus&& status, Yql::DqsProto::PingSessionResponse&&) mutable { - auto this_ = self.lock(); - if (!this_) { - return; - } - - if (status.GRpcStatusCode == grpc::INVALID_ARGUMENT || status.GRpcStatusCode == grpc::CANCELLED) { - YQL_CLOG(INFO, ProviderDq) << "Session closed " << sessionId; - this_->DropSession(sessionId); - } else { - this_->SchedulePingSessionRequest(sessionId); - } - }; - TaskScheduler.Delay(TDuration::Seconds(10)).Subscribe([self, callback, sessionId](const TFuture<void>&) { - auto this_ = self.lock(); - if (!this_) { - return; - } - - Yql::DqsProto::PingSessionRequest query; - query.SetSession(sessionId); - - this_->Service->DoRequest<Yql::DqsProto::PingSessionRequest, Yql::DqsProto::PingSessionResponse>( - query, - callback, - &Yql::DqsProto::DqService::Stub::AsyncPingSession); - }); - } - - TFuture<void> CloseSessionAsync(const TString& sessionId) { - std::shared_ptr<TDqGatewaySession> session; - with_lock (Mutex) { - auto it = Sessions.find(sessionId); - if (it != Sessions.end()) { - session = it->second; - Sessions.erase(it); - } - } - if (session) { - return session->Close(); - } - return MakeFuture(); - } - - TFuture<TResult> ExecutePlan(const TString& sessionId, NDqs::TPlan&& plan, const TVector<TString>& columns, - const THashMap<TString, TString>& secureParams, const THashMap<TString, TString>& graphParams, - const TDqSettings::TPtr& settings, - const TDqProgressWriter& progressWriter, const THashMap<TString, TString>& modulesMapping, - bool discard, ui64 executionTimeout) - { - std::shared_ptr<TDqGatewaySession> session; - with_lock(Mutex) { - auto it = Sessions.find(sessionId); - if (it != Sessions.end()) { - session = it->second; - } - } - TFailureInjector::Reach("dq_session_was_closed", [&] { session = nullptr; }); - if (!session) { - YQL_CLOG(ERROR, ProviderDq) << "Session was closed: " << sessionId; - auto res = NCommon::ResultFromException<TResult>(yexception() << "Session was closed"); - res.Fallback = true; - res.SetSuccess(); - return MakeFuture(res); - } - return session->ExecutePlan(std::move(plan), columns, secureParams, graphParams, settings, progressWriter, modulesMapping, discard, executionTimeout) - .Apply([](const TFuture<TResult>& f) { - try { - f.TryRethrow(); - } catch (const std::exception& e) { - YQL_CLOG(ERROR, ProviderDq) << e.what(); - return MakeFuture(NCommon::ResultFromException<TResult>(e)); - } - return f; - }); - } - -private: - NYdbGrpc::TGRpcClientConfig GrpcConf; - NYdbGrpc::TGRpcClientLow GrpcClient; - std::unique_ptr<NYdbGrpc::TServiceConnection<Yql::DqsProto::DqService>> Service; - - TDqTaskScheduler TaskScheduler; - const TDuration OpenSessionTimeout; - - TMutex Mutex; - THashMap<TString, std::shared_ptr<TDqGatewaySession>> Sessions; - - std::atomic<bool> IsStopped; -}; - -class TDqGateway: public IDqGateway { -public: - TDqGateway(const TString& host, int port, const TString& vanillaJobPath, const TString& vanillaJobMd5, TDuration timeout = TDuration::Minutes(60), TDuration requestTimeout = TDuration::Max()) - : Impl(std::make_shared<TDqGatewayImpl>(host, port, timeout, requestTimeout)) - , VanillaJobPath(vanillaJobPath) - , VanillaJobMd5(vanillaJobMd5) - { - } - - ~TDqGateway() { - Stop(); - } - - void Stop() override { - Impl->Stop(); - } - - TFuture<void> OpenSession(const TString& sessionId, const TString& username) override { - return Impl->OpenSession(sessionId, username); - } - - TFuture<void> CloseSessionAsync(const TString& sessionId) override { - return Impl->CloseSessionAsync(sessionId); - } - - TFuture<TResult> ExecutePlan(const TString& sessionId, NDqs::TPlan&& plan, const TVector<TString>& columns, - const THashMap<TString, TString>& secureParams, const THashMap<TString, TString>& graphParams, - const TDqSettings::TPtr& settings, - const TDqProgressWriter& progressWriter, const THashMap<TString, TString>& modulesMapping, - bool discard, ui64 executionTimeout) override - { - return Impl->ExecutePlan(sessionId, std::move(plan), columns, secureParams, graphParams, settings, progressWriter, modulesMapping, discard, executionTimeout); - } - - TString GetVanillaJobPath() override { - return VanillaJobPath; - } - - TString GetVanillaJobMd5() override { - return VanillaJobMd5; - } - -private: - std::shared_ptr<TDqGatewayImpl> Impl; - TString VanillaJobPath; - TString VanillaJobMd5; -}; - -TIntrusivePtr<IDqGateway> CreateDqGateway(const TString& host, int port) { - return new TDqGateway(host, port, "", ""); -} - -TIntrusivePtr<IDqGateway> CreateDqGateway(const NProto::TDqConfig& config) { - return new TDqGateway("localhost", config.GetPort(), - config.GetYtBackends()[0].GetVanillaJobLite(), - config.GetYtBackends()[0].GetVanillaJobLiteMd5(), - TDuration::MilliSeconds(config.GetOpenSessionTimeoutMs()), - TDuration::MilliSeconds(config.GetRequestTimeoutMs())); -} - -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/provider/yql_dq_gateway.h b/ydb/library/yql/providers/dq/provider/yql_dq_gateway.h index 4ea2fca8744..8ebc323646e 100644 --- a/ydb/library/yql/providers/dq/provider/yql_dq_gateway.h +++ b/ydb/library/yql/providers/dq/provider/yql_dq_gateway.h @@ -18,10 +18,6 @@ namespace NYql { -namespace NProto { -class TDqConfig; -} - class IDqGateway : public TThrRefBase { public: struct TStageStats { @@ -117,7 +113,4 @@ public: virtual void Stop() { } }; -TIntrusivePtr<IDqGateway> CreateDqGateway(const TString& host, int port); -TIntrusivePtr<IDqGateway> CreateDqGateway(const NProto::TDqConfig& config); - } // namespace NYql diff --git a/ydb/library/yql/providers/dq/provider/yql_dq_provider_ut.cpp b/ydb/library/yql/providers/dq/provider/yql_dq_provider_ut.cpp deleted file mode 100644 index 072e428f65b..00000000000 --- a/ydb/library/yql/providers/dq/provider/yql_dq_provider_ut.cpp +++ /dev/null @@ -1,365 +0,0 @@ -#include "yql_dq_statistics_json.h" - -#include <library/cpp/testing/unittest/registar.h> - -#include <yt/yql/providers/yt/gateway/file/yql_yt_file.h> -#include <yt/yql/providers/yt/gateway/file/yql_yt_file_services.h> -#include <yt/yql/providers/yt/lib/ut_common/yql_ut_common.h> -#include <yt/yql/providers/yt/provider/yql_yt_provider.h> - -#include <ydb/library/yql/dq/actors/compute/dq_compute_actor_async_io_factory.h> -#include <ydb/library/yql/dq/comp_nodes/yql_common_dq_factory.h> -#include <ydb/library/yql/dq/transform/yql_common_dq_transform.h> - -#include <yql/essentials/providers/common/comp_nodes/yql_factory.h> -#include <yql/essentials/providers/common/proto/gateways_config.pb.h> -#include <yql/essentials/providers/common/provider/yql_provider_names.h> - -#include <ydb/library/yql/providers/dq/common/yql_dq_common.h> -#include <ydb/library/yql/providers/dq/counters/counters.h> -#include <ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h> -#include <ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.h> -#include <ydb/library/yql/providers/dq/provider/yql_dq_gateway.h> -#include <ydb/library/yql/providers/dq/provider/yql_dq_provider.h> - -#include <yql/essentials/core/cbo/simple/cbo_simple.h> -#include <yql/essentials/core/facade/yql_facade.h> -#include <yql/essentials/core/file_storage/file_storage.h> -#include <yql/essentials/core/file_storage/proto/file_storage.pb.h> -#include <yql/essentials/core/services/mounts/yql_mounts.h> -#include <yql/essentials/minikql/comp_nodes/mkql_factories.h> -#include <yql/essentials/minikql/invoke_builtins/mkql_builtins.h> -#include <yql/essentials/minikql/mkql_function_registry.h> -#include <yql/essentials/utils/log/log.h> - -#include <util/stream/tee.h> -#include <util/string/cast.h> - -using namespace NYql; - -namespace { -// Runs a YQL/SQL program through the DQ engine with a YT file data source. -// maxTasksPerOperation sets the DQ limit on tasks per query. -bool RunDqProgram( - const TString& code, - ui32 maxTasksPerOperation, - const THashMap<TString, TString>& tableFiles, - TString* errorsMessage = nullptr) -{ - NLog::YqlLoggerScope logger("cerr", false); - - IOutputStream* errorsOutput = &Cerr; - TMaybe<TStringOutput> errorsMessageOutput; - TMaybe<TTeeOutput> tee; - if (errorsMessage) { - errorsMessageOutput.ConstructInPlace(*errorsMessage); - tee.ConstructInPlace(&*errorsMessageOutput, &Cerr); - errorsOutput = &*tee; - } - - TGatewaysConfig gatewaysConfig; - { - auto& dqCfg = *gatewaysConfig.MutableDq(); - auto addSetting = [&](const TString& name, const TString& value) { - auto* s = dqCfg.AddDefaultSettings(); - s->SetName(name); - s->SetValue(value); - }; - addSetting("EnableComputeActor", "1"); - addSetting("MaxTasksPerOperation", ToString(maxTasksPerOperation)); - } - - auto functionRegistry = NKikimr::NMiniKQL::CreateFunctionRegistry( - NKikimr::NMiniKQL::CreateBuiltinRegistry())->Clone(); - - TVector<TDataProviderInitializer> dataProvidersInit; - - auto yqlNativeServices = NFile::TYtFileServices::Make(functionRegistry.Get(), tableFiles); - auto ytGateway = CreateYtFileGateway(yqlNativeServices); - dataProvidersInit.push_back( - GetYtNativeDataProviderInitializer(ytGateway, MakeSimpleCBOOptimizerFactory(), {})); - - auto dqCompFactory = NKikimr::NMiniKQL::GetCompositeWithBuiltinFactory({ - NYql::GetCommonDqFactory(), - NKikimr::NMiniKQL::GetYqlFactory() - }); - auto dqTaskTransformFactory = NYql::CreateCompositeTaskTransformFactory({ - NYql::CreateCommonDqTaskTransformFactory() - }); - auto dqGateway = CreateLocalDqGateway( - functionRegistry.Get(), dqCompFactory, dqTaskTransformFactory, - {}, false, MakeIntrusive<NYql::NDq::TDqAsyncIoFactory>()); - - auto storage = NYql::CreateAsyncFileStorage({}); - dataProvidersInit.push_back( - NYql::GetDqDataProviderInitializer( - &CreateDqExecTransformer, dqGateway, dqCompFactory, {}, storage)); - - TExprContext moduleCtx; - IModuleResolver::TPtr moduleResolver; - YQL_ENSURE(GetYqlDefaultModuleResolver(moduleCtx, moduleResolver)); - - TProgramFactory factory(true, functionRegistry.Get(), 0ULL, dataProvidersInit, "ut"); - factory.SetGatewaysConfig(&gatewaysConfig); - factory.SetModules(moduleResolver); - - auto program = factory.Create("program", code); - - NSQLTranslation::TTranslationSettings sqlSettings; - sqlSettings.SyntaxVersion = 1; - sqlSettings.V0Behavior = NSQLTranslation::EV0Behavior::Disable; - sqlSettings.Flags.insert("DqEngineEnable"); - sqlSettings.Flags.insert("DqEngineForce"); - sqlSettings.ClusterMapping["plato"] = TString(YtProviderName); - - if (!program->ParseSql(sqlSettings)) { - program->PrintErrorsTo(*errorsOutput); - return false; - } - - if (!program->Compile("user")) { - program->PrintErrorsTo(*errorsOutput); - return false; - } - - auto status = program->Run("user", nullptr, nullptr, nullptr); - if (status == TProgram::TStatus::Error) { - program->PrintErrorsTo(*errorsOutput); - return false; - } - - return true; -} -} - -Y_UNIT_TEST_SUITE(TestCommon) { - -Y_UNIT_TEST(ParseCounterName) { - TString prefix = "Prefix"; - std::map<TString, TString> labels = { - {"Label1", "Value1"}, - {"Lavel2", "Value2"} - }; - TString name = "CounterName"; - - TString counterName = TCounters::GetCounterName(prefix, labels, name); - - TString prefix2, name2; - std::map<TString, TString> labels2; - - NCommon::ParseCounterName( - &prefix2, &labels2, &name2, counterName - ); - - UNIT_ASSERT_EQUAL(prefix, prefix2); - UNIT_ASSERT_EQUAL(name, name2); - UNIT_ASSERT_EQUAL(labels, labels2); -} - -Y_UNIT_TEST(CollectTaskRunnerStatisticsByStage) { - TOperationStatistics taskRunner; - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"Input", "2"}, - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter1"), - 1, 2, 3, 4, 5 - ) - ); - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"Output", "3"}, - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter2"), - 1, 2, 3, 4, 5 - ) - ); - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter3"), - 1, 2, 3, 4, 5 - ) - ); - - TStringStream result; - { - NYson::TYsonWriter writer(&result, NYT::NYson::EYsonFormat::Pretty); - CollectTaskRunnerStatisticsByStage(writer, taskRunner, true); - } - TString expected = R"__({ - "Stage=10" = { - "Input" = { - "total" = { - "Counter1" = { - "sum" = 1; - "count" = 5; - "avg" = 0; - "max" = 2; - "min" = 3 - } - } - }; - "Output" = { - "total" = { - "Counter2" = { - "sum" = 1; - "count" = 5; - "avg" = 0; - "max" = 2; - "min" = 3 - } - } - }; - "Source" = { - "total" = {} - }; - "Sink" = { - "total" = {} - }; - "Task" = { - "total" = { - "Counter3" = { - "sum" = 1; - "count" = 5; - "avg" = 0; - "max" = 2; - "min" = 3 - } - } - } - } -})__"; - UNIT_ASSERT_STRINGS_EQUAL(expected, result.Str()); -} - -Y_UNIT_TEST(CollectTaskRunnerStatisticsByTask) { - TOperationStatistics taskRunner; - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"InputChannel", "2"}, - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter1"), - 1, 2, 3, 4, 5 - ) - ); - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"OutputChannel", "3"}, - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter2"), - 1, 2, 3, 4, 5 - ) - ); - - taskRunner.Entries.push_back( - TOperationStatistics::TEntry( - TCounters::GetCounterName( - "Prefix", - { - {"Stage", "10"}, - {"Task", "1"}, - }, - "Counter3"), - 1, 2, 3, 4, 5 - ) - ); - - TStringStream result; - { - NYson::TYsonWriter writer(&result, NYT::NYson::EYsonFormat::Pretty); - CollectTaskRunnerStatisticsByTask(writer, taskRunner); - } - TString expected = R"__({ - "Stage=10" = { - "Task=1" = { - "Input=2" = { - "Counter1" = { - "sum" = 1; - "count" = 5; - "avg" = 4; - "max" = 2; - "min" = 3 - } - }; - "Output=3" = { - "Counter2" = { - "sum" = 1; - "count" = 5; - "avg" = 4; - "max" = 2; - "min" = 3 - } - }; - "Generic" = { - "Counter3" = { - "sum" = 1; - "count" = 5; - "avg" = 4; - "max" = 2; - "min" = 3 - } - } - } - } -})__"; - UNIT_ASSERT_STRINGS_EQUAL(expected, result.Str()); -} - -} - -Y_UNIT_TEST_SUITE(YqlDqExecTests) { - -// Reading from a YT table creates DQ source stages. The number of stages -// serves as a lower bound for the number of tasks (at least one task per -// stage). When stagesCount exceeds maxTasksPerOperation and no fallback is -// allowed (DqEngineForce), the exec transformer must report a user-friendly -// stages error message rather than crashing with YQL_ENSURE. -Y_UNIT_TEST(TooManyStagesErrorMessage) { - const TString code = R"( -USE plato; -SELECT key FROM Input; - )"; - - TTestTablesMapping testTables; - THashMap<TString, TString> tableFiles = { - {"yt.plato.Input", testTables.TmpInput.Name()}, - {"yt.plato.Output", testTables.TmpOutput.Name()}, - }; - - TString errorMessage; - bool ok = RunDqProgram(code, /*maxTasksPerOperation=*/0, tableFiles, &errorMessage); - UNIT_ASSERT_C(!ok, "Expected failure: too many stages"); - UNIT_ASSERT_STRING_CONTAINS(errorMessage, "stages exceeds the limit"); -} -} diff --git a/ydb/library/yql/providers/dq/service/grpc_service.cpp b/ydb/library/yql/providers/dq/service/grpc_service.cpp deleted file mode 100644 index e53fa64b211..00000000000 --- a/ydb/library/yql/providers/dq/service/grpc_service.cpp +++ /dev/null @@ -1,863 +0,0 @@ -#include "grpc_service.h" - -#include <yql/essentials/utils/log/log.h> - -#include <ydb/library/yql/providers/dq/actors/actor_helpers.h> -#include <ydb/library/yql/providers/dq/actors/executer_actor.h> -#include <ydb/library/yql/providers/dq/worker_manager/interface/events.h> -#include <ydb/library/yql/providers/dq/actors/execution_helpers.h> -#include <ydb/library/yql/providers/dq/actors/result_aggregator.h> -#include <ydb/library/yql/providers/dq/actors/events.h> -#include <ydb/library/yql/providers/dq/actors/task_controller.h> -#include <ydb/library/yql/providers/dq/actors/graph_execution_events_actor.h> - -#include <ydb/library/yql/providers/dq/counters/task_counters.h> -#include <ydb/library/yql/providers/dq/common/yql_dq_settings.h> -#include <ydb/library/yql/providers/dq/common/yql_dq_common.h> - -//#include <yql/tools/yqlworker/dq/worker_manager/benchmark.h> - -#include <yql/essentials/public/issue/yql_issue_message.h> - -#include <yql/essentials/minikql/invoke_builtins/mkql_builtins.h> - -#include <ydb/library/grpc/server/grpc_counters.h> -#include <ydb/public/api/protos/ydb_status_codes.pb.h> - -#include <ydb/library/actors/interconnect/interconnect.h> -#include <ydb/library/actors/core/subsystems/stats.h> -#include <ydb/library/actors/helpers/future_callback.h> -#include <library/cpp/build_info/build_info.h> -#include <library/cpp/svnversion/svnversion.h> - -#include <util/string/split.h> -#include <util/system/env.h> - -namespace NYql::NDqs { - using namespace NYql::NDqs; - using namespace NKikimr; - using namespace NThreading; - using namespace NMonitoring; - using namespace NActors; - - namespace { - NYdbGrpc::ICounterBlockPtr BuildCB(TIntrusivePtr<NMonitoring::TDynamicCounters>& counters, const TString& name) { - auto grpcCB = counters->GetSubgroup("rpc_name", name); - return MakeIntrusive<NYdbGrpc::TCounterBlock>( - grpcCB->GetCounter("total", true), - grpcCB->GetCounter("infly", true), - grpcCB->GetCounter("notOkReq", true), - grpcCB->GetCounter("notOkResp", true), - grpcCB->GetCounter("reqBytes", true), - grpcCB->GetCounter("inflyReqBytes", true), - grpcCB->GetCounter("resBytes", true), - grpcCB->GetCounter("notAuth", true), - grpcCB->GetCounter("resExh", true), - grpcCB); - } - - template<typename RequestType, typename ResponseType> - class TServiceProxyActor: public TSynchronizableRichActor<TServiceProxyActor<RequestType, ResponseType>> { - public: - static constexpr char ActorName[] = "SERVICE_PROXY"; - static constexpr char RetryName[] = "OperationRetry"; - - explicit TServiceProxyActor( - NYdbGrpc::IRequestContextBase* ctx, - const TIntrusivePtr<NMonitoring::TDynamicCounters>& counters, - const TString& traceId, const TString& username) - : TSynchronizableRichActor<TServiceProxyActor<RequestType, ResponseType>>(&TServiceProxyActor::Handler) - , Ctx(ctx) - , Counters(counters) - , ServiceProxyActorCounters(counters->GetSubgroup("component", "ServiceProxyActor")) - , ClientDisconnectedCounter(ServiceProxyActorCounters->GetCounter("ClientDisconnected", /*derivative=*/ true)) - , RetryCounter(ServiceProxyActorCounters->GetCounter(RetryName, /*derivative=*/ true)) - , FallbackCounter(ServiceProxyActorCounters->GetCounter("Fallback", /*derivative=*/ true)) - , ErrorCounter(ServiceProxyActorCounters->GetCounter("UnrecoverableError", /*derivative=*/ true)) - , Request(dynamic_cast<const RequestType*>(ctx->GetRequest())) - , TraceId(traceId) - , Username(username) - , Promise(NewPromise<void>()) - { - Settings->Dispatch(Request->GetSettings()); - Settings->FreezeDefaults(); - - MaxRetries = Settings->MaxRetries.Get().GetOrElse(MaxRetries); - RetryBackoff = TDuration::MilliSeconds(Settings->RetryBackoffMs.Get().GetOrElse(RetryBackoff.MilliSeconds())); - } - - STRICT_STFUNC(Handler, { - HFunc(TEvQueryResponse, OnReturnResult); - HFunc(TEvQueryStatus, OnQueryStatus); - cFunc(TEvents::TEvPoison::EventType, OnPoison); - SFunc(TEvents::TEvBootstrap, DoBootstrap); - hFunc(TEvDqStats, Handle); - }) - - TAutoPtr<IEventHandle> AfterRegister(const TActorId& self, const TActorId& parentId) override { - return new IEventHandle(self, parentId, new TEvents::TEvBootstrap(), 0); - } - - void OnPoison() { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(TraceId); - YQL_CLOG(DEBUG, ProviderDq) << __FUNCTION__ ; - ReplyError(grpc::UNAVAILABLE, "Unexpected error"); - *ClientDisconnectedCounter += 1; - } - - void Handle(TEvDqStats::TPtr&) { - // Do nothing - } - - void DoPassAway() override { - Promise.SetValue(); - } - - void DoBootstrap(const NActors::TActorContext& ctx) { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(TraceId); - if (!CtxSubscribed) { - auto selfId = ctx.SelfID; - auto* actorSystem = ctx.ActorSystem(); - Ctx->GetFinishFuture().Subscribe([selfId, actorSystem](const NYdbGrpc::IRequestContextBase::TAsyncFinishResult& future) { - Y_ABORT_UNLESS(future.HasValue()); - if (future.GetValue() == NYdbGrpc::IRequestContextBase::EFinishStatus::CANCEL) { - actorSystem->Send(selfId, new TEvents::TEvPoison()); - } - }); - CtxSubscribed = true; - } - Bootstrap(); - } - - virtual void Bootstrap() = 0; - - void SendResponse(TEvQueryResponse::TPtr& ev) - { - auto& result = ev->Get()->Record; - Yql::DqsProto::ExecuteQueryResult queryResult; - queryResult.Mutablesample()->CopyFrom(result.sample()); - - auto statusCode = result.GetStatusCode(); - // this code guarantees that query will be considered failed unless the status is SUCCESS - // fallback may be performed as an extra measure - if (statusCode != NYql::NDqProto::StatusIds::SUCCESS) { - YQL_CLOG(ERROR, ProviderDq) << "Query is considered FAILED, status=" << static_cast<int>(statusCode); - NYql::TIssue rootIssue("Fatal Error"); - rootIssue.SetCode(NCommon::NeedFallback(statusCode) ? TIssuesIds::DQ_GATEWAY_NEED_FALLBACK_ERROR : TIssuesIds::DQ_GATEWAY_ERROR, TSeverityIds::S_ERROR); - NYql::TIssues issues; - NYql::IssuesFromMessage(result.GetIssues(), issues); - for (const auto& issue: issues) { - rootIssue.AddSubIssue(MakeIntrusive<TIssue>(issue)); - } - result.MutableIssues()->Clear(); - NYql::IssuesToMessage({rootIssue}, result.MutableIssues()); - } - - for (const auto& [k, v] : QueryStat.Get()) { - std::map<TString, TString> labels; - TString prefix, name; - if (NCommon::ParseCounterName(&prefix, &labels, &name, k)) { - if (prefix == "Actor") { - auto group = Counters->GetSubgroup("counters", "Actor"); - for (const auto& [k, v] : labels) { - group = group->GetSubgroup(k, v); - } - group->GetHistogram(name, ExponentialHistogram(10, 2, 50000))->Collect(v.Sum); - } - } - } - - if (Settings->AggregateStatsByStage.Get().GetOrElse(TDqSettings::TDefault::AggregateStatsByStage)) { - auto aggregatedQueryStat = AggregateQueryStatsByStage(QueryStat, Task2Stage); - aggregatedQueryStat.FlushCounters(ResponseBuffer); - } else { - QueryStat.FlushCounters(ResponseBuffer); - } - - auto& operation = *ResponseBuffer.mutable_operation(); - operation.Setready(true); - operation.Mutableresult()->PackFrom(queryResult); - *operation.Mutableissues() = result.GetIssues(); - ResponseBuffer.SetTruncated(result.GetTruncated()); - ResponseBuffer.SetTimeout(result.GetTimeout()); - - Reply(Ydb::StatusIds::SUCCESS, statusCode > 1 || result.GetIssues().size() > 0); - } - - virtual void DoRetry() - { - this->CleanupChildren(); - auto selfId = this->SelfId(); - TActivationContext::Schedule(RetryBackoff, new IEventHandle(selfId, selfId, new TEvents::TEvBootstrap(), 0)); - Retry += 1; - *RetryCounter +=1 ; - } - - void OnQueryStatus(TEvQueryStatus::TPtr& ev, const TActorContext& ctx) { - Y_UNUSED(ev); Y_UNUSED(ctx); - if (!ExecuterActorId) { - auto response = MakeHolder<TEvQueryStatusResponse>(); - auto* r = response->Record.MutableResponse(); - r->SetStatus("Awaiting"); - this->Send(ev->Sender, response.Release()); - } else { - ctx.Send(ev->Forward(ExecuterActorId)); - } - } - - void OnReturnResult(TEvQueryResponse::TPtr& ev, const NActors::TActorContext& ctx) { - auto& result = ev->Get()->Record; - Y_UNUSED(ctx); - YQL_LOG_CTX_ROOT_SESSION_SCOPE(TraceId); - YQL_CLOG(DEBUG, ProviderDq) << "TServiceProxyActor::OnReturnResult " << result.GetMetric().size(); - QueryStat.AddCounters(result); - - auto statusCode = result.GetStatusCode(); - if ((statusCode != NYql::NDqProto::StatusIds::SUCCESS || result.GetIssues().size() > 0) && NCommon::IsRetriable(statusCode)) { - if (Retry < MaxRetries) { - QueryStat.AddCounter(RetryName, TDuration::MilliSeconds(0)); - NYql::TIssues issues; - NYql::IssuesFromMessage(result.GetIssues(), issues); - YQL_CLOG(WARN, ProviderDq) << RetryName << " " << Retry << " Issues: " << issues.ToString(); - DoRetry(); - } else { - YQL_CLOG(ERROR, ProviderDq) << "Retries limit exceeded, status= " << static_cast<int>(ev->Get()->Record.GetStatusCode()); - if (statusCode == NYql::NDqProto::StatusIds::SUCCESS) { - ev->Get()->Record.SetStatusCode(NYql::NDqProto::StatusIds::INTERNAL_ERROR); - } - SendResponse(ev); - } - } else { - SendResponse(ev); - } - } - - TFuture<void> GetFuture() { - return Promise.GetFuture(); - } - - virtual void ReplyError(grpc::StatusCode code, const TString& msg) { - Ctx->ReplyError(code, msg); - this->PassAway(); - } - - virtual void Reply(ui32 status, bool hasIssues) { - Y_UNUSED(hasIssues); - Ctx->Reply(&ResponseBuffer, status); - this->PassAway(); - } - - private: - NYdbGrpc::IRequestContextBase* Ctx; - bool CtxSubscribed = false; - ResponseType ResponseBuffer; - - protected: - TIntrusivePtr<NMonitoring::TDynamicCounters> Counters; - TIntrusivePtr<NMonitoring::TDynamicCounters> ServiceProxyActorCounters; - TDynamicCounters::TCounterPtr ClientDisconnectedCounter; - TDynamicCounters::TCounterPtr RetryCounter; - TDynamicCounters::TCounterPtr FallbackCounter; - TDynamicCounters::TCounterPtr ErrorCounter; - - const RequestType* Request; - const TString TraceId; - const TString Username; - TPromise<void> Promise; - const TInstant RequestStartTime = TInstant::Now(); - TActorId ExecuterActorId; - - TDqConfiguration::TPtr Settings = MakeIntrusive<TDqConfiguration>(); - - int Retry = 0; - int MaxRetries = 10; - TDuration RetryBackoff = TDuration::MilliSeconds(1000); - - NYql::TTaskCounters QueryStat; - THashMap<ui64, ui64> Task2Stage; - - void RestoreRequest() { - Request = dynamic_cast<const RequestType*>(Ctx->GetRequest()); - } - }; - - class TExecuteGraphProxyActor: public TServiceProxyActor<Yql::DqsProto::ExecuteGraphRequest, Yql::DqsProto::ExecuteGraphResponse> { - public: - using TBase = TServiceProxyActor<Yql::DqsProto::ExecuteGraphRequest, Yql::DqsProto::ExecuteGraphResponse>; - TExecuteGraphProxyActor(NYdbGrpc::IRequestContextBase* ctx, - const TIntrusivePtr<NMonitoring::TDynamicCounters>& counters, - const TString& traceId, const TString& username, - const NActors::TActorId& graphExecutionEventsActorId) - : TServiceProxyActor(ctx, counters, traceId, username) - , GraphExecutionEventsActorId(graphExecutionEventsActorId) - { - ExecutionTimeout = Request->GetExecutionTimeout(); - } - - void DoRetry() override { - YQL_CLOG(DEBUG, ProviderDq) << __FUNCTION__; - SendEvent(NYql::NDqProto::EGraphExecutionEventType::FAIL, nullptr, [this](const auto& ev) { - if (ev->Get()->Record.GetErrorMessage()) { - TBase::ReplyError(grpc::UNAVAILABLE, ev->Get()->Record.GetErrorMessage()); - } else { - RestoreRequest(); - ModifiedRequest.Reset(); - TBase::DoRetry(); - } - }); - } - - void Reply(ui32 status, bool hasIssues) override { - auto eventType = hasIssues - ? NYql::NDqProto::EGraphExecutionEventType::FAIL - : NYql::NDqProto::EGraphExecutionEventType::SUCCESS; - SendEvent(eventType, nullptr, [this, status, hasIssues](const auto& ev) { - if (!hasIssues && ev->Get()->Record.GetErrorMessage()) { - TBase::ReplyError(grpc::UNAVAILABLE, ev->Get()->Record.GetErrorMessage()); - } else { - TBase::Reply(status, hasIssues); - } - }); - } - - void ReplyError(grpc::StatusCode code, const TString& msg) override { - SendEvent(NYql::NDqProto::EGraphExecutionEventType::FAIL, nullptr, [this, code, msg](const auto& ev) { - Y_UNUSED(ev); - TBase::ReplyError(code, msg); - }); - } - - private: - THolder<Yql::DqsProto::ExecuteGraphRequest> ModifiedRequest; - - void DoPassAway() override { - YQL_CLOG(DEBUG, ProviderDq) << __FUNCTION__; - Send(GraphExecutionEventsActorId, new TEvents::TEvPoison()); - TServiceProxyActor::DoPassAway(); - } - - NDqProto::TGraphExecutionEvent::TExecuteGraphDescriptor SerializeGraphDescriptor() const { - NDqProto::TGraphExecutionEvent::TExecuteGraphDescriptor result; - - for (const auto& kv : Request->GetSecureParams()) { - result.MutableSecureParams()->MutableData()->insert(kv); - } - - for (const auto& kv : Request->GetGraphParams()) { - result.MutableGraphParams()->MutableData()->insert(kv); - } - - return result; - } - - void Bootstrap() override { - YQL_CLOG(DEBUG, ProviderDq) << "TServiceProxyActor::OnExecuteGraph"; - - SendEvent(NYql::NDqProto::EGraphExecutionEventType::START, SerializeGraphDescriptor(), [this](const TEvGraphExecutionEvent::TPtr& ev) { - if (ev->Get()->Record.GetErrorMessage()) { - TBase::ReplyError(grpc::UNAVAILABLE, ev->Get()->Record.GetErrorMessage()); - } else { - NDqProto::TGraphExecutionEvent::TMap params; - ev->Get()->Record.GetMessage().UnpackTo(¶ms); - FinishBootstrap(params); - } - }); - } - - void MergeTaskMetas(const NDqProto::TGraphExecutionEvent::TMap& params) { - if (!params.data().empty()) { - for (size_t i = 0; i < Request->TaskSize(); ++i) { - if (!ModifiedRequest) { - ModifiedRequest.Reset(new Yql::DqsProto::ExecuteGraphRequest()); - ModifiedRequest->CopyFrom(*Request); - } - - auto* task = ModifiedRequest->MutableTask(i); - - Yql::DqsProto::TTaskMeta taskMeta; - task->GetMeta().UnpackTo(&taskMeta); - - for (const auto&[key, value] : params.data()) { - (*taskMeta.MutableTaskParams())[key] = value; - } - - task->MutableMeta()->PackFrom(taskMeta); - } - } - - if (ModifiedRequest) { - Request = ModifiedRequest.Get(); - } - } - - void FinishBootstrap(const NDqProto::TGraphExecutionEvent::TMap& params) { - YQL_CLOG(DEBUG, ProviderDq) << __FUNCTION__; - MergeTaskMetas(params); - - ExecuterActorId = RegisterChild(NDq::MakeDqExecuter(MakeWorkerManagerActorID(SelfId().NodeId()), SelfId(), TraceId, Username, Settings, Counters, RequestStartTime, false, ExecutionTimeout)); - - TVector<TString> columns; - columns.reserve(Request->GetColumns().size()); - for (const auto& column : Request->GetColumns()) { - columns.push_back(column); - } - for (const auto& task : Request->GetTask()) { - Yql::DqsProto::TTaskMeta taskMeta; - task.GetMeta().UnpackTo(&taskMeta); - Task2Stage[task.GetId()] = taskMeta.GetStageId(); - } - THashMap<TString, TString> secureParams; - for (const auto& x : Request->GetSecureParams()) { - secureParams[x.first] = x.second; - } - auto resultId = RegisterChild(NExecutionHelpers::MakeResultAggregator( - columns, - ExecuterActorId, - TraceId, - secureParams, - Settings, - Request->GetResultType(), - Request->GetDiscard(), - GraphExecutionEventsActorId).Release()); - auto controlId = Settings->EnableComputeActor.Get().GetOrElse(false) == false ? resultId - : RegisterChild(NYql::MakeTaskController(TraceId, ExecuterActorId, resultId, NActors::TActorId{}, Settings, NYql::NCommon::TServiceCounters(Counters, nullptr, ""), TDuration::Seconds(5)).Release()); - Send(ExecuterActorId, MakeHolder<TEvGraphRequest>( - *Request, - controlId, - resultId)); - } - - template <class TPayload, class TCallback> - void SendEvent(NYql::NDqProto::EGraphExecutionEventType eventType, const TPayload& payload, TCallback callback) { - NDqProto::TGraphExecutionEvent record; - record.SetEventType(eventType); - if constexpr (!std::is_same_v<TPayload, std::nullptr_t>) { - record.MutableMessage()->PackFrom(payload); - } - Send(GraphExecutionEventsActorId, new TEvGraphExecutionEvent(record)); - Synchronize<TEvGraphExecutionEvent>([callback, traceId = TraceId](TEvGraphExecutionEvent::TPtr& ev) { - YQL_LOG_CTX_ROOT_SESSION_SCOPE(traceId); - Y_ABORT_UNLESS(ev->Get()->Record.GetEventType() == NYql::NDqProto::EGraphExecutionEventType::SYNC); - callback(ev); - }); - } - - NActors::TActorId GraphExecutionEventsActorId; - ui64 ExecutionTimeout; - }; - - TString GetVersionString() { - TStringBuilder sb; - sb << GetProgramSvnVersion() << "\n"; - sb << GetBuildInfo(); - TString sandboxTaskId = GetSandboxTaskId(); - if (sandboxTaskId != TString("0")) { - sb << "\nSandbox task id: " << sandboxTaskId; - } - - return sb; - } - } - - TDqsGrpcService::TDqsGrpcService( - NActors::TActorSystem& system, - TIntrusivePtr<NMonitoring::TDynamicCounters> counters, - const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories) - : ActorSystem(system) - , Counters(std::move(counters)) - , DqTaskPreprocessorFactories(dqTaskPreprocessorFactories) - , Promise(NewPromise<void>()) - , RunningRequests(0) - , Stopping(false) - , Sessions(&system, Counters->GetSubgroup("component", "grpc")->GetCounter("Sessions")) - { } - -#define ADD_REQUEST(NAME, IN, OUT, ACTION) \ - do { \ - MakeIntrusive<NYdbGrpc::TGRpcRequest<Yql::DqsProto::IN, Yql::DqsProto::OUT, TDqsGrpcService>>( \ - this, \ - &Service_, \ - CQ, \ - [this](NYdbGrpc::IRequestContextBase* ctx) { ACTION }, \ - &Yql::DqsProto::DqService::AsyncService::Request##NAME, \ - #NAME, \ - logger, \ - BuildCB(Counters, #NAME)) \ - ->Run(); \ - } while (0) - - void TDqsGrpcService::InitService(grpc::ServerCompletionQueue* cq, NYdbGrpc::TLoggerPtr logger) { - using namespace google::protobuf; - - CQ = cq; - - using TDqTaskPreprocessorCollection = std::vector<NYql::IDqTaskPreprocessor::TPtr>; - - ADD_REQUEST(ExecuteGraph, ExecuteGraphRequest, ExecuteGraphResponse, { - TGuard<TMutex> lock(Mutex); - if (Stopping) { - ctx->ReplyError(grpc::UNAVAILABLE, "Terminating in progress"); - return; - } - auto* request = dynamic_cast<const Yql::DqsProto::ExecuteGraphRequest*>(ctx->GetRequest()); - auto session = Sessions.GetSession(request->GetSession()); - uint64_t querySeqNo = request->GetQuerySeqNo(); - if (!session) { - TString message = TStringBuilder() - << "Bad session: " - << request->GetSession(); - YQL_CLOG(DEBUG, ProviderDq) << message; - ctx->ReplyError(grpc::INVALID_ARGUMENT, message); - return; - } - - TDqTaskPreprocessorCollection taskPreprocessors; - for (const auto& factory: DqTaskPreprocessorFactories) { - taskPreprocessors.push_back(factory()); - } - - auto graphExecutionEventsActorId = ActorSystem.Register(NDqs::MakeGraphExecutionEventsActor(request->GetSession(), std::move(taskPreprocessors))); - - RunningRequests++; - auto actor = MakeHolder<TExecuteGraphProxyActor>(ctx, Counters, request->GetSession(), session->GetUsername(), graphExecutionEventsActorId); - auto future = actor->GetFuture(); - auto actorId = ActorSystem.Register(actor.Release()); - future.Apply([session, actorId, this] (const TFuture<void>&) mutable { - RunningRequests--; - if (Stopping && !RunningRequests) { - Promise.SetValue(); - } - - session->DeleteRequest(actorId); - }); - session->AddRequest(actorId, querySeqNo); - }); - - ADD_REQUEST(SvnRevision, SvnRevisionRequest, SvnRevisionResponse, { - Y_UNUSED(this); - Yql::DqsProto::SvnRevisionResponse result; - result.SetRevision(GetVersionString()); - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - }); - - ADD_REQUEST(CloseSession, CloseSessionRequest, CloseSessionResponse, { - Y_UNUSED(this); - auto* request = dynamic_cast<const Yql::DqsProto::CloseSessionRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - Yql::DqsProto::CloseSessionResponse result; - Sessions.CloseSession(request->GetSession()); - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - }); - - ADD_REQUEST(PingSession, PingSessionRequest, PingSessionResponse, { - Y_UNUSED(this); - auto* request = dynamic_cast<const Yql::DqsProto::PingSessionRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - YQL_CLOG(TRACE, ProviderDq) << "PingSession " << request->GetSession(); - - Yql::DqsProto::PingSessionResponse result; - auto session = Sessions.GetSession(request->GetSession()); - if (!session) { - TString message = TStringBuilder() - << "Bad session: " - << request->GetSession(); - YQL_CLOG(DEBUG, ProviderDq) << message; - ctx->ReplyError(grpc::INVALID_ARGUMENT, message); - } else { - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - } - }); - - ADD_REQUEST(OpenSession, OpenSessionRequest, OpenSessionResponse, { - Y_UNUSED(this); - auto* request = dynamic_cast<const Yql::DqsProto::OpenSessionRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - YQL_CLOG(DEBUG, ProviderDq) << "OpenSession for " << request->GetSession() << " " << request->GetUsername(); - - Yql::DqsProto::OpenSessionResponse result; - if (Sessions.OpenSession(request->GetSession(), request->GetUsername())) { - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - } else { - ctx->ReplyError(grpc::INVALID_ARGUMENT, "Session `" + request->GetSession() + "' exists"); - } - }); - - ADD_REQUEST(JobStop, JobStopRequest, JobStopResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::JobStopRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - auto ev = MakeHolder<TEvJobStop>(*request); - - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::JobStopResponse>(ctx->GetArena()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - - ActorSystem.Send(MakeWorkerManagerActorID(ActorSystem.NodeId), ev.Release()); - }); - - ADD_REQUEST(ClusterStatus, ClusterStatusRequest, ClusterStatusResponse, { - auto ev = MakeHolder<TEvClusterStatus>(); - - using ResultEv = TEvClusterStatusResponse; - - TExecutorPoolStats poolStats; - - TExecutorThreadStats stat; - TVector<TExecutorThreadStats> stats; - GetActorSystemStats(ActorSystem).GetPoolStats(0, poolStats, stats); - - for (const auto& s : stats) { - stat.Aggregate(s); - } - - YQL_CLOG(DEBUG, ProviderDq) << "SentEvents: " << stat.SentEvents; - YQL_CLOG(DEBUG, ProviderDq) << "ReceivedEvents: " << stat.ReceivedEvents; - YQL_CLOG(DEBUG, ProviderDq) << "NonDeliveredEvents: " << stat.NonDeliveredEvents; - YQL_CLOG(DEBUG, ProviderDq) << "EmptyMailboxActivation: " << stat.EmptyMailboxActivation; - Sessions.PrintInfo(); - - for (ui32 i = 0; i < stat.ActorsAliveByActivity.size(); i=i+1) { - if (stat.ActorsAliveByActivity[i]) { - YQL_CLOG(DEBUG, ProviderDq) << "ActorsAliveByActivity[" << i << "]=" << stat.ActorsAliveByActivity[i]; - } - } - - auto callback = MakeHolder<TRichActorFutureCallback<ResultEv>>( - [ctx] (TAutoPtr<TEventHandle<ResultEv>>& event) mutable { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::ClusterStatusResponse>(ctx->GetArena()); - result->MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "ClusterStatus failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(2000)); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, ev.Release(), IEventHandle::FlagTrackDelivery)); - }); - - ADD_REQUEST(OperationStop, OperationStopRequest, OperationStopResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::OperationStopRequest*>(ctx->GetRequest()); - auto ev = MakeHolder<TEvOperationStop>(*request); - - auto callback = MakeHolder<TRichActorFutureCallback<TEvOperationStopResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvOperationStopResponse>>& event) mutable { - Y_UNUSED(event); - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::OperationStopResponse>(ctx->GetArena()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "OperationStopResponse failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(2000)); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, ev.Release(), IEventHandle::FlagTrackDelivery)); - }); - - ADD_REQUEST(QueryStatus, QueryStatusRequest, QueryStatusResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::QueryStatusRequest*>(ctx->GetRequest()); - - auto ev = MakeHolder<TEvQueryStatus>(*request); - - auto callback = MakeHolder<TRichActorFutureCallback<TEvQueryStatusResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvQueryStatusResponse>>& event) mutable { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::QueryStatusResponse>(ctx->GetArena()); - result->MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "QueryStatus failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(2000)); - - uint64_t querySeqNo = request->GetQuerySeqNo(); - if (querySeqNo) { - auto session = Sessions.GetSession(request->GetSession()); - if (!session) { - TString message = TStringBuilder() - << "Bad session: " - << request->GetSession(); - YQL_CLOG(DEBUG, ProviderDq) << message; - ctx->ReplyError(grpc::INVALID_ARGUMENT, message); - } else { - auto actorId = session->FindActorId(querySeqNo); - if (!actorId) { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::QueryStatusResponse>(ctx->GetArena()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - } else { - TActorId callbackId = ActorSystem.Register(callback.Release()); - ActorSystem.Send(new IEventHandle(actorId, callbackId, ev.Release())); - } - } - } else { - TActorId callbackId = ActorSystem.Register(callback.Release()); - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, ev.Release(), IEventHandle::FlagTrackDelivery)); - } - }); - - ADD_REQUEST(RegisterNode, RegisterNodeRequest, RegisterNodeResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::RegisterNodeRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - if (!request->GetPort() - || request->GetRole().empty() - || request->GetAddress().empty()) - { - ctx->ReplyError(grpc::INVALID_ARGUMENT, "Invalid argument"); - return; - } - - auto ev = MakeHolder<TEvRegisterNode>(*request); - - using ResultEv = TEvRegisterNodeResponse; - - auto callback = MakeHolder<TRichActorFutureCallback<ResultEv>>( - [ctx] (TAutoPtr<TEventHandle<ResultEv>>& event) mutable { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::RegisterNodeResponse>(ctx->GetArena()); - result->MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "RegisterNode failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(5000)); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, ev.Release(), IEventHandle::FlagTrackDelivery)); - }); - - ADD_REQUEST(GetMaster, GetMasterRequest, GetMasterResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::GetMasterRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - auto requestEvent = MakeHolder<TEvGetMasterRequest>(); - - auto callback = MakeHolder<TActorFutureCallback<TEvGetMasterResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvGetMasterResponse>>& event) mutable { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::GetMasterResponse>(ctx->GetArena()); - result->MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, requestEvent.Release())); - }); - - ADD_REQUEST(ConfigureFailureInjector, ConfigureFailureInjectorRequest, ConfigureFailureInjectorResponse,{ - auto* request = dynamic_cast<const Yql::DqsProto::ConfigureFailureInjectorRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - auto requestEvent = MakeHolder<TEvConfigureFailureInjectorRequest>(*request); - - auto callback = MakeHolder<TActorFutureCallback<TEvConfigureFailureInjectorResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvConfigureFailureInjectorResponse>>& event) mutable { - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::ConfigureFailureInjectorResponse>(ctx->GetArena()); - result->MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, requestEvent.Release())); - }); - - ADD_REQUEST(IsReady, IsReadyRequest, IsReadyResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::IsReadyRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - auto ev = MakeHolder<TEvIsReady>(*request); - - auto callback = MakeHolder<TRichActorFutureCallback<TEvIsReadyResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvIsReadyResponse>>& event) mutable { - Yql::DqsProto::IsReadyResponse result; - result.SetIsReady(event->Get()->Record.GetIsReady()); - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "IsReadyForRevision failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(2000)); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(ActorSystem.NodeId), callbackId, ev.Release())); - }); - - ADD_REQUEST(Routes, RoutesRequest, RoutesResponse, { - auto* request = dynamic_cast<const Yql::DqsProto::RoutesRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!request); - - auto ev = MakeHolder<TEvRoutesRequest>(); - - auto callback = MakeHolder<TRichActorFutureCallback<TEvRoutesResponse>>( - [ctx] (TAutoPtr<TEventHandle<TEvRoutesResponse>>& event) mutable { - Yql::DqsProto::RoutesResponse result; - result.MergeFrom(event->Get()->Record.GetResponse()); - ctx->Reply(&result, Ydb::StatusIds::SUCCESS); - }, - [ctx] () mutable { - YQL_CLOG(INFO, ProviderDq) << "Routes failed"; - ctx->ReplyError(grpc::UNAVAILABLE, "Error"); - }, - TDuration::MilliSeconds(5000)); - - TActorId callbackId = ActorSystem.Register(callback.Release()); - - ActorSystem.Send(new IEventHandle(MakeWorkerManagerActorID(request->GetNodeId()), callbackId, ev.Release())); - }); - -/* 1. move grpc to providers/dq, 2. move benchmark to providers/dq 3. uncomment - ADD_REQUEST(Benchmark, BenchmarkRequest, BenchmarkResponse, { - auto* req = dynamic_cast<const Yql::DqsProto::BenchmarkRequest*>(ctx->GetRequest()); - Y_ABORT_UNLESS(!!req); - - TWorkerManagerBenchmarkOptions options; - if (req->GetWorkerCount()) { - options.WorkerCount = req->GetWorkerCount(); - } - if (req->GetInflight()) { - options.Inflight = req->GetInflight(); - } - if (req->GetTotalRequests()) { - options.TotalRequests = req->GetTotalRequests(); - } - if (req->GetMaxRunTimeMs()) { - options.MaxRunTimeMs = TDuration::MilliSeconds(req->GetMaxRunTimeMs()); - } - - auto benchmarkId = ActorSystem.Register( - CreateWorkerManagerBenchmark( - MakeWorkerManagerActorID(ActorSystem.NodeId), options - )); - - ActorSystem.Send(benchmarkId, new TEvents::TEvBootstrap); - - auto* result = google::protobuf::Arena::CreateMessage<Yql::DqsProto::BenchmarkResponse>(ctx->GetArena()); - ctx->Reply(result, Ydb::StatusIds::SUCCESS); - }); -*/ - } - - TFuture<void> TDqsGrpcService::Stop() { - TGuard<TMutex> lock(Mutex); - Stopping = true; - auto future = Promise.GetFuture(); - if (RunningRequests == 0) { - Promise.SetValue(); - } - return future; - } -} diff --git a/ydb/library/yql/providers/dq/service/grpc_service.h b/ydb/library/yql/providers/dq/service/grpc_service.h deleted file mode 100644 index c0479f3a7d4..00000000000 --- a/ydb/library/yql/providers/dq/service/grpc_service.h +++ /dev/null @@ -1,47 +0,0 @@ -#pragma once - -#include <ydb/library/yql/providers/dq/interface/yql_dq_task_preprocessor.h> - -#include <ydb/library/yql/providers/dq/api/grpc/api.grpc.pb.h> -#include <ydb/library/yql/providers/dq/api/protos/service.pb.h> - -#include <yql/essentials/minikql/mkql_function_registry.h> - -#include <ydb/library/grpc/server/grpc_request.h> -#include <ydb/library/grpc/server/grpc_server.h> - -#include <ydb/library/actors/core/actorsystem.h> -#include <ydb/library/actors/core/event_local.h> -#include <ydb/library/actors/core/events.h> -#include <library/cpp/monlib/dynamic_counters/counters.h> -#include <library/cpp/threading/future/future.h> - -#include "grpc_session.h" - -namespace NYql::NDqs { - class TDatabaseManager; - - class TDqsGrpcService: public NYdbGrpc::TGrpcServiceBase<Yql::DqsProto::DqService> { - public: - TDqsGrpcService(NActors::TActorSystem& system, - TIntrusivePtr<NMonitoring::TDynamicCounters> counters, - const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories); - - void InitService(grpc::ServerCompletionQueue* cq, NYdbGrpc::TLoggerPtr logger) override; - - NThreading::TFuture<void> Stop(); - - private: - NActors::TActorSystem& ActorSystem; - grpc::ServerCompletionQueue* CQ = nullptr; - - TIntrusivePtr<NMonitoring::TDynamicCounters> Counters; - TDqTaskPreprocessorFactoryCollection DqTaskPreprocessorFactories; - TMutex Mutex; - NThreading::TPromise<void> Promise; - std::atomic<ui64> RunningRequests; - std::atomic<bool> Stopping; - - TSessionStorage Sessions; - }; -} diff --git a/ydb/library/yql/providers/dq/service/grpc_session.cpp b/ydb/library/yql/providers/dq/service/grpc_session.cpp deleted file mode 100644 index a523af899b6..00000000000 --- a/ydb/library/yql/providers/dq/service/grpc_session.cpp +++ /dev/null @@ -1,124 +0,0 @@ -#include "grpc_session.h" - -#include <yql/essentials/utils/log/log.h> - -namespace NYql::NDqs { - -TSession::~TSession() { - TGuard<TMutex> lock(Mutex); - for (auto [actorId, _] : ActorId2QueryNo) { - ActorSystem->Send(actorId, new NActors::TEvents::TEvPoison()); - } -} - -void TSession::DeleteRequest(const NActors::TActorId& actorId) -{ - TGuard<TMutex> lock(Mutex); - auto it = ActorId2QueryNo.find(actorId); - if (it != ActorId2QueryNo.end()) { - if (it->second) { - QueryNo2ActorId.erase(it->second); - } - ActorId2QueryNo.erase(it); - } -} - -void TSession::AddRequest(const NActors::TActorId& actorId, uint64_t querySeqNo) -{ - TGuard<TMutex> lock(Mutex); - ActorId2QueryNo[actorId] = querySeqNo; - if (querySeqNo) { - QueryNo2ActorId[querySeqNo] = actorId; - } -} - -NActors::TActorId TSession::FindActorId(uint64_t queryNo) { - TGuard<TMutex> lock(Mutex); - auto it = QueryNo2ActorId.find(queryNo); - if (it != QueryNo2ActorId.end()) { - return it->second; - } - return {}; -} - -TSessionStorage::TSessionStorage( - NActors::TActorSystem* actorSystem, - const NMonitoring::TDynamicCounters::TCounterPtr& sessionsCounter) - : ActorSystem(actorSystem) - , SessionsCounter(sessionsCounter) -{ } - -void TSessionStorage::CloseSession(const TString& sessionId) -{ - TGuard<TMutex> lock(SessionMutex); - auto it = Sessions.find(sessionId); - if (it == Sessions.end()) { - return; - } - SessionsByLastUpdate.erase(it->second.Iterator); - Sessions.erase(it); - *SessionsCounter = Sessions.size(); -} - -std::shared_ptr<TSession> TSessionStorage::GetSession(const TString& sessionId) -{ - Clean(TInstant::Now() - TDuration::Minutes(10)); - - TGuard<TMutex> lock(SessionMutex); - auto it = Sessions.find(sessionId); - if (it == Sessions.end()) { - return std::shared_ptr<TSession>(); - } else { - SessionsByLastUpdate.erase(it->second.Iterator); - SessionsByLastUpdate.push_back({TInstant::Now(), sessionId}); - it->second.Iterator = SessionsByLastUpdate.end(); - it->second.Iterator--; - return it->second.Session; - } -} - -bool TSessionStorage::OpenSession(const TString& sessionId, const TString& username) -{ - TGuard<TMutex> lock(SessionMutex); - if (Sessions.contains(sessionId)) { - return false; - } - - SessionsByLastUpdate.push_back({TInstant::Now(), sessionId}); - auto it = SessionsByLastUpdate.end(); --it; - - Sessions[sessionId] = TSessionAndIterator { - std::make_shared<TSession>(username, ActorSystem), - it - }; - - *SessionsCounter = Sessions.size(); - - return true; -} - -void TSessionStorage::Clean(TInstant before) { - TGuard<TMutex> lock(SessionMutex); - for (TSessionsByLastUpdate::iterator it = SessionsByLastUpdate.begin(); - it != SessionsByLastUpdate.end(); ) - { - if (it->LastUpdate < before) { - YQL_CLOG(INFO, ProviderDq) << "Drop session by timeout " << it->SessionId; - Sessions.erase(it->SessionId); - it = SessionsByLastUpdate.erase(it); - } else { - break; - } - } - - *SessionsCounter = Sessions.size(); -} - -void TSessionStorage::PrintInfo() const { - YQL_CLOG(INFO, ProviderDq) << "SessionsByLastUpdate: " << SessionsByLastUpdate.size(); - YQL_CLOG(DEBUG, ProviderDq) << "Sessions: " << Sessions.size(); - ui64 currenSessionsCounter = *SessionsCounter; - YQL_CLOG(DEBUG, ProviderDq) << "SessionsCounter: " << currenSessionsCounter; -} - -} // namespace NYql::NDqs diff --git a/ydb/library/yql/providers/dq/service/grpc_session.h b/ydb/library/yql/providers/dq/service/grpc_session.h deleted file mode 100644 index e71ad3dc512..00000000000 --- a/ydb/library/yql/providers/dq/service/grpc_session.h +++ /dev/null @@ -1,59 +0,0 @@ -#pragma once -#include <ydb/library/actors/core/actorsystem.h> - -namespace NYql::NDqs { - -class TSession { -public: - TSession(const TString& username, NActors::TActorSystem* actorSystem) - : Username(username) - , ActorSystem(actorSystem) - { } - - ~TSession(); - - const TString& GetUsername() const { - return Username; - } - - void DeleteRequest(const NActors::TActorId& id); - void AddRequest(const NActors::TActorId& id, uint64_t queryNo); - NActors::TActorId FindActorId(uint64_t queryNo); - -private: - TMutex Mutex; - const TString Username; - NActors::TActorSystem* ActorSystem; - THashMap<NActors::TActorId, uint64_t> ActorId2QueryNo; - THashMap<uint64_t, NActors::TActorId> QueryNo2ActorId; -}; - -class TSessionStorage { -public: - TSessionStorage( - NActors::TActorSystem* actorSystem, - const NMonitoring::TDynamicCounters::TCounterPtr& sessionsCounter); - void CloseSession(const TString& sessionId); - std::shared_ptr<TSession> GetSession(const TString& sessionId); - bool OpenSession(const TString& sessionId, const TString& username); - void Clean(TInstant before); - void PrintInfo() const; - -private: - NActors::TActorSystem* ActorSystem; - NMonitoring::TDynamicCounters::TCounterPtr SessionsCounter; - TMutex SessionMutex; - struct TTimeAndSessionId { - TInstant LastUpdate; - TString SessionId; - }; - using TSessionsByLastUpdate = TList<TTimeAndSessionId>; - TSessionsByLastUpdate SessionsByLastUpdate; - struct TSessionAndIterator { - std::shared_ptr<TSession> Session; - TSessionsByLastUpdate::iterator Iterator; - }; - THashMap<TString, TSessionAndIterator> Sessions; -}; - -} // namespace NYql::NDqs diff --git a/ydb/library/yql/providers/dq/service/interconnect_helpers.cpp b/ydb/library/yql/providers/dq/service/interconnect_helpers.cpp deleted file mode 100644 index e872a511e28..00000000000 --- a/ydb/library/yql/providers/dq/service/interconnect_helpers.cpp +++ /dev/null @@ -1,317 +0,0 @@ -#include "interconnect_helpers.h" -#include "service_node.h" - -#include "grpc_service.h" - -#include <ydb/library/actors/helpers/selfping_actor.h> - -#include <yql/essentials/utils/log/log.h> -#include <yql/essentials/utils/backtrace/backtrace.h> -#include <yql/essentials/utils/yql_panic.h> - -#include <yql/essentials/minikql/invoke_builtins/mkql_builtins.h> - -#include <ydb/library/actors/core/executor_pool_basic.h> -#include <ydb/library/actors/core/scheduler_basic.h> -#include <ydb/library/actors/core/scheduler_actor.h> -#include <ydb/library/actors/dnsresolver/dnsresolver.h> -#include <ydb/library/actors/interconnect/interconnect.h> -#include <ydb/library/actors/interconnect/interconnect_common.h> -#include <ydb/library/actors/interconnect/interconnect_tcp_proxy.h> -#include <ydb/library/actors/interconnect/interconnect_tcp_server.h> -#include <ydb/library/actors/interconnect/poller/poller_actor.h> -#include <library/cpp/yson/node/node_io.h> - -#include <util/stream/file.h> -#include <util/system/env.h> -#include <util/system/fs.h> - -namespace NYql::NDqs { - using namespace NActors; - using namespace NActors::NDnsResolver; - using namespace NYdbGrpc; - - class TYqlLogBackend: public TLogBackend { - void WriteData(const TLogRecord& rec) override { - TString message(rec.Data, rec.Len); - if (message.find("ICP01 ready to work") != TString::npos) { - return; - } - YQL_CLOG(DEBUG, ProviderDq) << message; - } - - void ReopenLog() override { } - }; - - static void InitSelfPingActor(NActors::TActorSystemSetup* setup, NMonitoring::TDynamicCounterPtr rootCounters) - { - const TDuration selfPingInterval = TDuration::MilliSeconds(10); - - const auto counters = rootCounters->GetSubgroup("counters", "utils"); - - for (size_t poolId = 0; poolId < setup->GetExecutorsCount(); ++poolId) { - const auto& poolName = setup->GetPoolName(poolId); - auto poolGroup = counters->GetSubgroup("execpool", poolName); - auto selfPinfMaxCounter = poolGroup->GetCounter("SelfPingMaxUs", false); - auto selfPinfAvgCounter = poolGroup->GetCounter("SelfPingAvgUs", false); - auto selfPinfAvgCounterIn1s = poolGroup->GetCounter("SelfPingAvgUsIn1s", false); - auto cpuTimeCounter = poolGroup->GetCounter("CpuMatBenchNs", false); - IActor* selfPingActor = CreateSelfPingActor(selfPingInterval, selfPinfMaxCounter, - selfPinfAvgCounter, selfPinfAvgCounterIn1s, cpuTimeCounter); - setup->LocalServices.push_back( - std::make_pair(TActorId(), - TActorSetupCmd(selfPingActor, - TMailboxType::HTSwap, - poolId))); - } - } - - std::tuple<THolder<NActors::TActorSystemSetup>, TIntrusivePtr<NActors::NLog::TSettings>> BuildActorSetup( - ui32 nodeId, - TString interconnectAddress, - ui16 port, - SOCKET socket, - TVector<ui32> threads, - NMonitoring::TDynamicCounterPtr counters, - const TNameserverFactory& nameserverFactory, - TMaybe<ui32> maxNodeId, - const NYql::NProto::TDqConfig::TICSettings& icSettings) - { - auto setup = MakeHolder<TActorSystemSetup>(); - - setup->NodeId = nodeId; - - if (threads.empty()) { - threads = {icSettings.GetThreads()}; - } - - setup->ExecutorsCount = threads.size(); - setup->Executors.Reset(new TAutoPtr<IExecutorPool>[setup->ExecutorsCount]); - for (ui32 i = 0; i < setup->ExecutorsCount; ++i) { - setup->Executors[i] = new TBasicExecutorPool( - i, - threads[i], - 50, - "pool-"+ToString(i)// poolName - ); - } - auto schedulerConfig = TSchedulerConfig(); - schedulerConfig.MonCounters = counters; - -#define SET_VALUE(name) \ - if (icSettings.Has ## name()) { \ - schedulerConfig.name = icSettings.Get ## name (); \ - YQL_CLOG(DEBUG, ProviderDq) << "Scheduler IC " << #name << " set to " << schedulerConfig.name; \ - } - - SET_VALUE(ResolutionMicroseconds); - SET_VALUE(SpinThreshold); - SET_VALUE(ProgressThreshold); - SET_VALUE(UseSchedulerActor); - SET_VALUE(RelaxedSendPaceEventsPerSecond); - SET_VALUE(RelaxedSendPaceEventsPerCycle); - SET_VALUE(RelaxedSendThresholdEventsPerSecond); - SET_VALUE(RelaxedSendThresholdEventsPerCycle); - -#undef SET_VALUE - - setup->Scheduler = CreateSchedulerThread(schedulerConfig); - - YQL_CLOG(DEBUG, ProviderDq) << "Initializing local services"; - setup->LocalServices.emplace_back(MakePollerActorId(), TActorSetupCmd(CreatePollerActor(), TMailboxType::ReadAsFilled, 0)); - if (IActor* schedulerActor = CreateSchedulerActor(schedulerConfig)) { - TActorId schedulerActorId = MakeSchedulerActorId(); - setup->LocalServices.emplace_back(schedulerActorId, TActorSetupCmd(schedulerActor, TMailboxType::ReadAsFilled, 0)); - } - - NActors::TActorId loggerActorId(nodeId, "logger"); - auto logSettings = MakeIntrusive<NActors::NLog::TSettings>(loggerActorId, - 0, NActors::NLog::PRI_INFO); - static TString defaultComponent = "ActorLib"; - logSettings->Append(0, 1024, [&](NActors::NLog::EComponent) -> const TString & { return defaultComponent; }); - TString explanation = ""; - if (YQL_CVLOG_ACTIVE(NLog::ELevel::TRACE, NLog::EComponent::CoreDq)) { - logSettings->SetLevel(NActors::NLog::PRI_TRACE, 535 /*NKikimrServices::KQP_COMPUTE*/, explanation); - logSettings->SetLevel(NActors::NLog::PRI_TRACE, 713 /*NKikimrServices::YQL_PROXY*/, explanation); - logSettings->SetLevel(NActors::NLog::PRI_TRACE, 1165 /*NKikimrServices::DQ_TASK_RUNNER*/, explanation); - } - NActors::TLoggerActor *loggerActor = new NActors::TLoggerActor( - logSettings, - new TYqlLogBackend, - counters->GetSubgroup("logger", "counters")); - setup->LocalServices.emplace_back(logSettings->LoggerActorId, TActorSetupCmd(loggerActor, TMailboxType::Simple, 0)); - - TIntrusivePtr<TTableNameserverSetup> nameserverTable = new TTableNameserverSetup(); - THashSet<ui32> staticNodeId; - - YQL_CLOG(DEBUG, ProviderDq) << "Initializing node table"; - nameserverTable->StaticNodeTable[nodeId] = std::make_pair(interconnectAddress, port); - - setup->LocalServices.emplace_back( - MakeDnsResolverActorId(), TActorSetupCmd(CreateOnDemandDnsResolver(), TMailboxType::ReadAsFilled, 0)); - - setup->LocalServices.emplace_back( - GetNameserviceActorId(), TActorSetupCmd(nameserverFactory(nameserverTable), TMailboxType::ReadAsFilled, 0)); - - - InitSelfPingActor(setup.Get(), counters); - - TIntrusivePtr<TInterconnectProxyCommon> icCommon = new TInterconnectProxyCommon(); - icCommon->NameserviceId = GetNameserviceActorId(); - Y_UNUSED(counters); - //icCommon->MonCounters = counters->GetSubgroup("counters", "interconnect"); - icCommon->MonCounters = MakeIntrusive<NMonitoring::TDynamicCounters>(); - -#define SET_DURATION(name) \ - { \ - icCommon->Settings.name = TDuration::MilliSeconds(icSettings.Get ## name ## Ms()); \ - YQL_CLOG(DEBUG, ProviderDq) << "IC " << #name << " set to " << icCommon->Settings.name; \ - } - -#define SET_VALUE(name) \ - { \ - icCommon->Settings.name = icSettings.Get ## name(); \ - YQL_CLOG(DEBUG, ProviderDq) << "IC " << #name << " set to " << icCommon->Settings.name; \ - } - - SET_DURATION(Handshake); - SET_DURATION(DeadPeer); - SET_DURATION(CloseOnIdle); - - SET_VALUE(SendBufferDieLimitInMB); - SET_VALUE(TotalInflightAmountOfData); - SET_VALUE(MergePerPeerCounters); - SET_VALUE(MergePerDataCenterCounters); - SET_VALUE(TCPSocketBufferSize); - - SET_DURATION(PingPeriod); - SET_DURATION(ForceConfirmPeriod); - SET_DURATION(LostConnection); - SET_DURATION(BatchPeriod); - - SET_DURATION(MessagePendingTimeout); - - SET_VALUE(MessagePendingSize); - SET_VALUE(MaxSerializedEventSize); - SET_VALUE(EnableExternalDataChannel); - -#undef SET_DURATION -#undef SET_VALUE - - YQL_CLOG(DEBUG, ProviderDq) << "Initializing proxy actors"; - auto effectiveMaxNodeId = maxNodeId.GetOrElse(static_cast<ui32>(ENodeIdLimits::MaxWorkerNodeId)); - setup->Interconnect.ProxyActors.resize(effectiveMaxNodeId + 1); - for (ui32 id = 1; id <= effectiveMaxNodeId; ++id) { - if (nodeId != id) { - IActor* actor = new TInterconnectProxyTCP(id, icCommon); - setup->Interconnect.ProxyActors[id] = TActorSetupCmd(actor, TMailboxType::ReadAsFilled, 0); - } - } - - // start listener - YQL_CLOG(DEBUG, ProviderDq) << "Start listener"; - { - icCommon->TechnicalSelfHostName = interconnectAddress; - YQL_CLOG(INFO, ProviderDq) << "Start listener " << interconnectAddress << ":" << port << " socket: " << socket; - IActor* listener; - TMaybe<SOCKET> maybeSocket = socket < 0 - ? Nothing() - : TMaybe<SOCKET>(socket); - - listener = new NActors::TInterconnectListenerTCP(interconnectAddress, port, icCommon, maybeSocket); - - setup->LocalServices.emplace_back( - MakeInterconnectListenerActorId(false), - TActorSetupCmd(listener, TMailboxType::ReadAsFilled, 0)); - } - - YQL_CLOG(DEBUG, ProviderDq) << "Actor initialization complete"; - -#ifdef _unix_ - signal(SIGPIPE, SIG_IGN); -#endif - - return std::make_tuple(std::move(setup), logSettings); - } - - std::tuple<TString, TString> GetLocalAddress(const TString* overrideHostname, int family) { - constexpr auto MaxLocalHostNameLength = 4096; - std::array<char, MaxLocalHostNameLength> buffer; - buffer.fill(0); - TString hostName; - TString localAddress; - - int result = gethostname(buffer.data(), buffer.size() - 1); - if (result != 0) { - Cerr << "gethostname failed for " << std::string_view{buffer.data(), buffer.size()} << " error " << strerror(errno) << Endl; - return std::make_tuple(hostName, localAddress); - } - - if (overrideHostname) { - memcpy(&buffer[0], overrideHostname->c_str(), Min<int>( - overrideHostname->size()+1, buffer.size()-1 - )); - } - - hostName = &buffer[0]; - - addrinfo request; - memset(&request, 0, sizeof(request)); - request.ai_family = family; - request.ai_socktype = SOCK_STREAM; - - addrinfo* response = nullptr; - result = getaddrinfo(buffer.data(), nullptr, &request, &response); - if (result != 0) { - Cerr << "getaddrinfo failed for " << std::string_view{buffer.data(), buffer.size()} << " error " << gai_strerror(result) << Endl; - return std::make_tuple(hostName, localAddress); - } - - std::unique_ptr<addrinfo, void (*)(addrinfo*)> holder(response, &freeaddrinfo); - - if (!response->ai_addr) { - Cerr << "getaddrinfo failed: no ai_addr" << Endl; - return std::make_tuple(hostName, localAddress); - } - - auto* sa = response->ai_addr; - Y_ABORT_UNLESS(sa->sa_family == family); - switch (family) { - case AF_INET6: - inet_ntop(AF_INET6, &(((struct sockaddr_in6*)sa)->sin6_addr), - &buffer[0], buffer.size() - 1); - break; - case AF_INET: - inet_ntop(AF_INET, &(((struct sockaddr_in*)sa)->sin_addr), - &buffer[0], buffer.size() - 1); - break; - default: - Y_ABORT_UNLESS(false); - break; - } - - localAddress = &buffer[0]; - - return std::make_tuple(hostName, localAddress); - } - - std::tuple<TString, TString> GetUserToken(const TMaybe<TString>& maybeUser, const TMaybe<TString>& maybeTokenFile) - { - auto home = GetEnv("HOME"); - auto systemUser = GetEnv("USER"); - - TString userName = maybeUser - ? *maybeUser - : systemUser; - - TString tokenFile = maybeTokenFile - ? *maybeTokenFile - : home + "/.yt/token"; - - TString token = NFs::Exists(tokenFile) - ? TFileInput(tokenFile).ReadLine() - : ""; - - return std::make_tuple(userName, token); - } -} diff --git a/ydb/library/yql/providers/dq/service/interconnect_helpers.h b/ydb/library/yql/providers/dq/service/interconnect_helpers.h deleted file mode 100644 index 73a38686b30..00000000000 --- a/ydb/library/yql/providers/dq/service/interconnect_helpers.h +++ /dev/null @@ -1,51 +0,0 @@ -#pragma once - -#include <ydb/library/actors/core/actorsystem.h> -#include <ydb/library/actors/interconnect/poller/poller_tcp.h> -#include <ydb/library/actors/interconnect/interconnect.h> -#include <library/cpp/yson/node/node.h> - -#include <ydb/library/yql/providers/dq/config/config.pb.h> - -namespace NYql::NDqs { - -enum class ENodeIdLimits { - MinServiceNodeId = 1, - MaxServiceNodeId = 512, // excluding - MinWorkerNodeId = 512, - MaxWorkerNodeId = 8192, // excluding -}; - -using TNameserverFactory = std::function<NActors::IActor*(const TIntrusivePtr<NActors::TTableNameserverSetup>& setup)>; - -struct TServiceNodeConfig { - ui32 NodeId; - TString InterconnectAddress; - TString GrpcHostname; - ui16 Port; - ui16 GrpcPort = 8080; - ui16 MbusPort = 0; - SOCKET Socket = -1; - SOCKET GrpcSocket = -1; - TMaybe<ui32> MaxNodeId; - NYql::NProto::TDqConfig::TICSettings ICSettings = NYql::NProto::TDqConfig::TICSettings(); - TNameserverFactory NameserverFactory = [](const TIntrusivePtr<NActors::TTableNameserverSetup>& setup) { - return CreateNameserverTable(setup); - }; -}; - -std::tuple<TString, TString> GetLocalAddress(const TString* hostname = nullptr, int family = AF_INET6); -std::tuple<TString, TString> GetUserToken(const TMaybe<TString>& user, const TMaybe<TString>& tokenFile); - -std::tuple<THolder<NActors::TActorSystemSetup>, TIntrusivePtr<NActors::NLog::TSettings>> BuildActorSetup( - ui32 nodeId, - TString interconnectAddress, - ui16 port, - SOCKET socket, - TVector<ui32> threads, - NMonitoring::TDynamicCounterPtr counters, - const TNameserverFactory& nameserverFactory, - TMaybe<ui32> maxNodeId, - const NYql::NProto::TDqConfig::TICSettings& icSettings = NYql::NProto::TDqConfig::TICSettings()); - -} // namespace NYql::NDqs diff --git a/ydb/library/yql/providers/dq/service/service_node.cpp b/ydb/library/yql/providers/dq/service/service_node.cpp deleted file mode 100644 index 6f82dec398a..00000000000 --- a/ydb/library/yql/providers/dq/service/service_node.cpp +++ /dev/null @@ -1,167 +0,0 @@ -#include "service_node.h" - -#include <yql/essentials/utils/log/log.h> -#include <yql/essentials/providers/common/metrics/metrics_registry.h> -#include <yql/essentials/utils/yql_panic.h> - -#include <ydb/library/yql/providers/dq/actors/execution_helpers.h> - -#include <ydb/library/grpc/server/actors/logger.h> - -#include <utility> - -namespace NYql { - using namespace NActors; - using namespace NYdbGrpc; - using namespace NYql::NDqs; - - class TGrpcExternalListener: public IExternalListener { - public: - TGrpcExternalListener(SOCKET socket) - : Socket(socket) - , Listener(MakeIntrusive<NInterconnect::TStreamSocket>(Socket)) - { - SetNonBlock(socket, true); - } - - void Init(std::unique_ptr<grpc::experimental::ExternalConnectionAcceptor> acceptor) override { - Acceptor = std::move(acceptor); - } - - private: - void Start() override { - YQL_CLOG(DEBUG, ProviderDq) << "Start GRPC Listener"; - Poller.Start(); - StartRead(); - } - - void Stop() override { - Poller.Stop(); - } - - void StartRead() { - YQL_CLOG(TRACE, ProviderDq) << "Read next GRPC event"; - Poller.StartRead(Listener, [&](const TIntrusivePtr<TSharedDescriptor>& ss) { - return Accept(ss); - }); - } - - TDelegate Accept(const TIntrusivePtr<TSharedDescriptor>& ss) - { - NInterconnect::TStreamSocket* socket = (NInterconnect::TStreamSocket*)ss.Get(); - int r = 0; - while (r >= 0) { - NInterconnect::TAddress address; - r = socket->Accept(address); - if (r >= 0) { - YQL_CLOG(TRACE, ProviderDq) << "New GRPC connection"; - grpc::experimental::ExternalConnectionAcceptor::NewConnectionParameters params; - SetNonBlock(r, true); - params.listener_fd = static_cast<int>(Socket); - params.fd = r; - Acceptor->HandleNewConnection(¶ms); - } else if (-r != EAGAIN && -r != EWOULDBLOCK) { - YQL_CLOG(DEBUG, ProviderDq) << "Unknown error code " + ToString(r); - } - } - - return [this] { - StartRead(); - }; - } - - SOCKET Socket; - TIntrusivePtr<NInterconnect::TStreamSocket> Listener; - NInterconnect::TPollerThreads Poller; - std::unique_ptr<grpc::experimental::ExternalConnectionAcceptor> Acceptor; - }; - - TServiceNode::TServiceNode( - const TServiceNodeConfig& config, - ui32 threads, - IMetricsRegistryPtr metricsRegistry) - : Config(config) - , Threads(threads) - , MetricsRegistry(std::move(metricsRegistry)) - { - std::tie(Setup, LogSettings) = BuildActorSetup( - Config.NodeId, - Config.InterconnectAddress, - Config.Port, - Config.Socket, - {Threads, 8}, - MetricsRegistry->GetSensors(), - Config.NameserverFactory, - Config.MaxNodeId, - Config.ICSettings); - } - - void TServiceNode::AddLocalService(TActorId actorId, TActorSetupCmd service) { - YQL_ENSURE(!ActorSystem); - Setup->LocalServices.emplace_back(actorId, std::move(service)); - } - - NActors::TActorSystem* TServiceNode::StartActorSystem(void* appData) { - Y_ABORT_UNLESS(!ActorSystem); - - ActorSystem = MakeHolder<NActors::TActorSystem>(Setup, appData, LogSettings); - ActorSystem->Start(); - - return ActorSystem.Get(); - } - - void TServiceNode::StartService(const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories) { - class TCustomOption : public grpc::ServerBuilderOption { - public: - TCustomOption() { } - - void UpdateArguments(grpc::ChannelArguments *args) override { - args->SetInt(GRPC_ARG_ALLOW_REUSEPORT, 1); - } - - void UpdatePlugins(std::vector<std::unique_ptr<grpc::ServerBuilderPlugin>>* /*plugins*/) override - { } - }; - - YQL_CLOG(INFO, ProviderDq) << "Starting GRPC on " << Config.GrpcPort; - - IExternalListener::TPtr listener = nullptr; - if (Config.GrpcSocket >= 0) { - listener = MakeIntrusive<TGrpcExternalListener>(Config.GrpcSocket); - } - - auto options = TServerOptions() - // .SetHost(CurrentNode->Address) - .SetHost("[::]") - .SetPort(Config.GrpcPort) - .SetExternalListener(listener) - .SetWorkerThreads(2) - .SetGRpcMemoryQuotaBytes(1024 * 1024 * 1024) - .SetMaxMessageSize(1024 * 1024 * 256) - .SetMaxGlobalRequestInFlight(50000) - .SetUseAuth(false) - .SetKeepAliveEnable(true) - .SetKeepAliveIdleTimeoutTriggerSec(360) - .SetKeepAliveMaxProbeCount(3) - .SetKeepAliveProbeIntervalSec(1) - .SetServerBuilderMutator([](grpc::ServerBuilder& builder) { - builder.SetOption(std::make_unique<TCustomOption>()); - }) - .SetLogger(CreateActorSystemLogger(*ActorSystem, 413)); // 413 - NKikimrServices::GRPC_SERVER - - Server = MakeHolder<TGRpcServer>(options); - Service = TIntrusivePtr<IGRpcService>(new TDqsGrpcService(*ActorSystem, MetricsRegistry->GetSensors(), dqTaskPreprocessorFactories)); - Server->AddService(Service); - Server->Start(); - } - - void TServiceNode::Stop(TDuration timeout) { - (static_cast<TDqsGrpcService*>(Service.Get()))->Stop().Wait(timeout); - - Server->Stop(); - for (auto id : ActorIds) { - ActorSystem->Send(id, new NActors::TEvents::TEvPoison); - } - ActorSystem->Stop(); - } -} // namespace NYql diff --git a/ydb/library/yql/providers/dq/service/service_node.h b/ydb/library/yql/providers/dq/service/service_node.h deleted file mode 100644 index 443a3bfbc2d..00000000000 --- a/ydb/library/yql/providers/dq/service/service_node.h +++ /dev/null @@ -1,48 +0,0 @@ -#pragma once - -#include "grpc_service.h" -#include "interconnect_helpers.h" - -#include <yql/essentials/providers/common/metrics/metrics_registry.h> -#include <ydb/library/yql/providers/dq/interface/yql_dq_task_preprocessor.h> - -#include <yql/essentials/minikql/mkql_function_registry.h> - -#include <ydb/library/actors/core/executor_pool_basic.h> -#include <ydb/library/actors/core/scheduler_basic.h> -#include <ydb/library/actors/interconnect/interconnect.h> -#include <ydb/library/actors/interconnect/interconnect_common.h> -#include <ydb/library/actors/interconnect/interconnect_tcp_proxy.h> -#include <ydb/library/actors/interconnect/interconnect_tcp_server.h> -#include <ydb/library/actors/interconnect/poller/poller_actor.h> - -namespace NYql { - class TServiceNode { - public: - TServiceNode( - const NDqs::TServiceNodeConfig& config, - ui32 threads, - IMetricsRegistryPtr metricsRegistry); - - void AddLocalService(NActors::TActorId actorId, NActors::TActorSetupCmd service); - NActors::TActorSystem* StartActorSystem(void* appData = nullptr); - void StartService(const TDqTaskPreprocessorFactoryCollection& dqTaskPreprocessorFactories); - - void Stop(TDuration time = TDuration::Max()); - - NActors::TActorSystemSetup* GetSetup() const { - return Setup.Get(); - } - - private: - NDqs::TServiceNodeConfig Config; - ui32 Threads; - IMetricsRegistryPtr MetricsRegistry; - THolder<NActors::TActorSystemSetup> Setup; - TIntrusivePtr<NActors::NLog::TSettings> LogSettings; - THolder<NActors::TActorSystem> ActorSystem; - TVector<NActors::TActorId> ActorIds; - THolder<NYdbGrpc::TGRpcServer> Server; - TIntrusivePtr<NYdbGrpc::IGRpcService> Service; - }; -} diff --git a/ydb/library/yql/providers/dq/service/ya.make b/ydb/library/yql/providers/dq/service/ya.make deleted file mode 100644 index 53f66ea0578..00000000000 --- a/ydb/library/yql/providers/dq/service/ya.make +++ /dev/null @@ -1,33 +0,0 @@ -LIBRARY() - -SRCS( - grpc_service.cpp - grpc_session.cpp - service_node.cpp - interconnect_helpers.cpp -) - -PEERDIR( - ydb/library/actors/core - ydb/library/actors/dnsresolver - ydb/library/actors/interconnect - library/cpp/build_info - ydb/library/grpc/server - ydb/library/grpc/server/actors - library/cpp/svnversion - library/cpp/threading/future - yql/essentials/sql - ydb/public/api/protos - yql/essentials/providers/common/metrics - ydb/library/yql/providers/dq/actors - ydb/library/yql/providers/dq/api/grpc - ydb/library/yql/providers/dq/common - ydb/library/yql/providers/dq/counters - ydb/library/yql/providers/dq/interface - ydb/library/yql/providers/dq/worker_manager - ydb/library/yql/providers/dq/worker_manager/interface -) - -YQL_LAST_ABI_VERSION() - -END() diff --git a/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.cpp b/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.cpp deleted file mode 100644 index dfb950e1d46..00000000000 --- a/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.cpp +++ /dev/null @@ -1,31 +0,0 @@ -#include "pool_stats_collector.h" - -#include <ydb/library/actors/helpers/pool_stats_collector.h> - -namespace NYql { - -using namespace NActors; -using namespace NMonitoring; - -namespace { - -TIntrusivePtr<TDynamicCounters> GetServiceCounters(TIntrusivePtr<TDynamicCounters> root, - const TString &service) -{ - auto res = root->GetSubgroup("counters", service); - auto utils = root->GetSubgroup("counters", "utils"); - auto lookupCounter = utils->GetSubgroup("component", service)->GetCounter("CounterLookups", true); - res->SetLookupCounter(lookupCounter); - return res; -} - -} - -IActor *CreateStatsCollector(ui32 intervalSec, - const TActorSystemSetup& setup, - NMonitoring::TDynamicCounterPtr counters) -{ - return new NActors::TStatsCollectingActor(intervalSec, setup, GetServiceCounters(counters, "utils")); -} - -} // namespace NKikimr diff --git a/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.h b/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.h deleted file mode 100644 index b29f16889a3..00000000000 --- a/ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.h +++ /dev/null @@ -1,18 +0,0 @@ -#pragma once - -#include <ydb/library/actors/core/actorsystem.h> -#include <library/cpp/monlib/dynamic_counters/counters.h> - -namespace NActors { - struct TActorSystemSetup; -} - -// copy from kikimr/core/base/pool_stats_collector.h - -namespace NYql { - - NActors::IActor* CreateStatsCollector(ui32 intervalSec, - const NActors::TActorSystemSetup& setup, - NMonitoring::TDynamicCounterPtr counters); - -} diff --git a/ydb/library/yql/providers/dq/stats_collector/ya.make b/ydb/library/yql/providers/dq/stats_collector/ya.make deleted file mode 100644 index c9fed95d42d..00000000000 --- a/ydb/library/yql/providers/dq/stats_collector/ya.make +++ /dev/null @@ -1,20 +0,0 @@ -LIBRARY() - -SET( - SOURCE - pool_stats_collector.cpp -) - -SRCS( - ${SOURCE} -) - -PEERDIR( - ydb/library/actors/core - ydb/library/actors/helpers - library/cpp/monlib/dynamic_counters -) - -YQL_LAST_ABI_VERSION() - -END() diff --git a/ydb/library/yql/providers/dq/ya.make b/ydb/library/yql/providers/dq/ya.make index 7d9818e9944..ca83e435bb0 100644 --- a/ydb/library/yql/providers/dq/ya.make +++ b/ydb/library/yql/providers/dq/ya.make @@ -10,8 +10,6 @@ RECURSE( planner provider runtime - service - stats_collector task_runner task_runner_actor worker_manager diff --git a/ydb/library/yql/tools/dqrun/.gitignore b/ydb/library/yql/tools/dqrun/.gitignore deleted file mode 100644 index 234301d156b..00000000000 --- a/ydb/library/yql/tools/dqrun/.gitignore +++ /dev/null @@ -1 +0,0 @@ -dummy_op.* diff --git a/ydb/library/yql/tools/dqrun/README.md b/ydb/library/yql/tools/dqrun/README.md deleted file mode 100644 index 528a92b6d2e..00000000000 --- a/ydb/library/yql/tools/dqrun/README.md +++ /dev/null @@ -1,132 +0,0 @@ -# dqrun - Utility for Local Debugging of Distributed SQL Engine - -`dqrun` is a utility designed for local debugging of a distributed SQL engine. It allows you to run all components of a distributed engine in a single process for more convenient debugging. - -## Command-Line Options - -- `-s`: if specified, SQL is used; if not specified, the query plan execution specified in s-expression is used. -- `-p <file>`: specify a file with an SQL query or a s-expression. -- `--gateways-cfg <file>`: specify a file with the engine configuration. (Example: [examples/gateways.conf](examples/gateways.conf)) -- `--fs-cfg <file>`: specify a file with the file cache configuration. (Example: [examples/fs.conf](examples/fs.conf)) -- `--bindings-file <file>`: specify a file with the data schema. (Examples: [examples/bindings_tpch.json](examples/bindings_tpch.json), [examples/bindings_tpch_pg.json](examples/bindings_tpch_pg.json)) -- `--dq-host <host>`: set the host to connect to the test utilities `service_node` and `worker_node`. -- `--dq-port <port>`: set the port to connect to the test utilities `service_node` and `worker_node`. - -## Example of Local Usage - -```bash -dqrun -s -p query.sql --gateways-cfg examples/gateways.conf --fs-cfg examples/fs.conf --bindings-file examples/bindings_tpch.json -``` - -In this example, `dqrun` will use SQL from the file `query.sql`, engine configuration from the file [examples/gateways.conf](examples/gateways.conf), file cache configuration from the file [examples/fs.conf](examples/fs.conf), and data schema from the file [examples/bindings_tpch.json](examples/bindings_tpch.json). - -Download `tpc` data using one of `download_files_h_*.sh` scripts from `yql/queries/tpc_benchmark` number in the name corresponds to the size of data. To download minimal working example use `download_files_h_1.sh`. Move the downloaded `tpc` folder to current directory. - -To run dq you will also need a query. The simple example of a query: - -```sql -SELECT * FROM bindings.customer LIMIT 10; -``` - -## Example of Usage as a Client to Test Utilities - -```bash -dqrun --dq-host localhost --dq-port 8080 -s -p query.sql --gateways-cfg examples/gateways.conf --fs-cfg examples/fs.conf --bindings-file examples/bindings_tpch.json -``` - -In this example, `dqrun` will use SQL from the file `query.sql`, engine configuration from the file [examples/gateways.conf](examples/gateways.conf), file cache configuration from the file [examples/fs.conf](examples/fs.conf), and data schema from the file [examples/bindings_tpch.json](examples/bindings_tpch.json). Additionally, the utility will act as a client to the test utilities `service_node` and `worker_node`, using the specified host and port. - -## Data Schema File Example (bindings_tpch.json) - -An example of a data schema file for queries from the TPC-H benchmark. - -```json -{ - "customer": { - "ClusterType": "s3", - "path": "h/1/customer/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", - [ - ["c_acctbal", ["DataType", "Double"]], - ["c_address", ["DataType", "String"]], - ["c_comment", ["DataType", "String"]], - ["c_custkey", ["DataType", "Int32"]], - ["c_mktsegment", ["DataType", "String"]], - ["c_name", ["DataType", "String"]], - ["c_nationkey", ["DataType", "Int32"]], - ["c_phone", ["DataType", "String"]] - ] - ] - }, - // Other tables from the TPC-H benchmark -} -``` - -## Data Schema File Example with PostgreSQL Syntax (bindings_tpch_pg.json) - -An example of a data schema file for queries from the TPC-H benchmark using PostgreSQL syntax. - -```json -{ - "customer": { - "ClusterType": "s3", - "path": "h/1/customer/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", - [ - ["c_acctbal", ["PgType", "numeric"]], - ["c_address", ["PgType", "text"]], - ["c_comment", ["PgType", "text"]], - ["c_custkey", ["PgType", "int4"]], - ["c_mktsegment", ["PgType", "text"]], - ["c_name", ["PgType", "text"]], - ["c_nationkey", ["PgType", "int4"]], - ["c_phone", ["PgType", "text"]] - ] - ] - }, - // Other tables from the TPC-H benchmark -} -``` - -**Example of gateways.conf:** - -```conf -Dq { - DefaultSettings { - Name: "HashJoinMode" - Value: "grace" - } - - DefaultSettings { - Name: "UseOOBTransport" - Value: "true" - } - - DefaultSettings { - Name: "UseWideChannels" - Value: "true" - } -} - -S3 { - ClusterMapping { - Name: "yq-clickbench-local" - Url: "file://./clickbench/" - } - ClusterMapping { - Name: "yq-tpc-local" - Url: "file://./tpc/" - } -} -``` - -In the `Dq` section, parameters for the distributed engine are specified. The complete list of parameters can be found [here](../../providers/dq/common/yql_dq_settings.h). - -In the `S3` section, either S3 clusters or directories pretending to be S3 clusters are specified. In this example, the "cluster" `yq-tpc-local` points to the `tpc` directory located locally. - diff --git a/ydb/library/yql/tools/dqrun/dqrun.cpp b/ydb/library/yql/tools/dqrun/dqrun.cpp deleted file mode 100644 index d21643f0f58..00000000000 --- a/ydb/library/yql/tools/dqrun/dqrun.cpp +++ /dev/null @@ -1,12 +0,0 @@ -#include <ydb/library/yql/tools/dqrun/lib/dqrun_lib.h> -#include <util/generic/yexception.h> - -int main(int argc, const char *argv[]) { - try { - return NYql::TDqRunTool().Main(argc, argv); - } - catch (...) { - Cerr << CurrentExceptionMessage() << Endl; - return 1; - } -} diff --git a/ydb/library/yql/tools/dqrun/examples/bindings_tpcds.json b/ydb/library/yql/tools/dqrun/examples/bindings_tpcds.json deleted file mode 100644 index 9bf2a9d466c..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/bindings_tpcds.json +++ /dev/null @@ -1,668 +0,0 @@ -{ - "call_center": { - "ClusterType": "s3", - "path": "ds/1/call_center/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["cc_call_center_id", ["OptionalType", ["DataType", "String"]]], - ["cc_call_center_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cc_city", ["OptionalType", ["DataType", "String"]]], - ["cc_class", ["OptionalType", ["DataType", "String"]]], - ["cc_closed_date_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cc_company", ["OptionalType", ["DataType", "Int32"]]], - ["cc_company_name", ["OptionalType", ["DataType", "String"]]], - ["cc_country", ["OptionalType", ["DataType", "String"]]], - ["cc_county", ["OptionalType", ["DataType", "String"]]], - ["cc_division", ["OptionalType", ["DataType", "Int32"]]], - ["cc_division_name", ["OptionalType", ["DataType", "String"]]], - ["cc_employees", ["OptionalType", ["DataType", "Int32"]]], - ["cc_gmt_offset", ["OptionalType", ["DataType", "Double"]]], - ["cc_hours", ["OptionalType", ["DataType", "String"]]], - ["cc_manager", ["OptionalType", ["DataType", "String"]]], - ["cc_market_manager", ["OptionalType", ["DataType", "String"]]], - ["cc_mkt_class", ["OptionalType", ["DataType", "String"]]], - ["cc_mkt_desc", ["OptionalType", ["DataType", "String"]]], - ["cc_mkt_id", ["OptionalType", ["DataType", "Int32"]]], - ["cc_name", ["OptionalType", ["DataType", "String"]]], - ["cc_open_date_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cc_rec_end_date", ["OptionalType", ["DataType", "String"]]], - ["cc_rec_start_date", ["OptionalType", ["DataType", "String"]]], - ["cc_sq_ft", ["OptionalType", ["DataType", "Int32"]]], - ["cc_state", ["OptionalType", ["DataType", "String"]]], - ["cc_street_name", ["OptionalType", ["DataType", "String"]]], - ["cc_street_number", ["OptionalType", ["DataType", "String"]]], - ["cc_street_type", ["OptionalType", ["DataType", "String"]]], - ["cc_suite_number", ["OptionalType", ["DataType", "String"]]], - ["cc_tax_percentage", ["OptionalType", ["DataType", "Double"]]], - ["cc_zip", ["OptionalType", ["DataType", "String"]]] - ] - ] - }, - "catalog_page": { - "ClusterType": "s3", - "path": "ds/1/catalog_page/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["cp_catalog_number", ["OptionalType", ["DataType", "Int32"]]], - ["cp_catalog_page_id", ["OptionalType", ["DataType", "String"]]], - ["cp_catalog_page_number", ["OptionalType", ["DataType", "Int32"]]], - ["cp_catalog_page_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cp_department", ["OptionalType", ["DataType", "String"]]], - ["cp_description", ["OptionalType", ["DataType", "String"]]], - ["cp_end_date_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cp_start_date_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cp_type", ["OptionalType", ["DataType", "String"]]] - ] - ] - }, - "catalog_returns": { - "ClusterType": "s3", - "path": "ds/1/catalog_returns/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["cr_call_center_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_catalog_page_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_fee", ["OptionalType", ["DataType", "Double"]]], - ["cr_item_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_net_loss", ["OptionalType", ["DataType", "Double"]]], - ["cr_order_number", ["OptionalType", ["DataType", "Int32"]]], - ["cr_reason_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_refunded_addr_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_refunded_cash", ["OptionalType", ["DataType", "Double"]]], - ["cr_refunded_cdemo_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_refunded_customer_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_refunded_hdemo_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_return_amount", ["OptionalType", ["DataType", "Double"]]], - ["cr_return_amt_inc_tax", ["OptionalType", ["DataType", "Double"]]], - ["cr_return_quantity", ["OptionalType", ["DataType", "Int32"]]], - ["cr_return_ship_cost", ["OptionalType", ["DataType", "Double"]]], - ["cr_return_tax", ["OptionalType", ["DataType", "Double"]]], - ["cr_returned_date_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_returned_time_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_returning_addr_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_returning_cdemo_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_returning_customer_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_returning_hdemo_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_reversed_charge", ["OptionalType", ["DataType", "Double"]]], - ["cr_ship_mode_sk", ["OptionalType", ["DataType", "Int32"]]], - ["cr_store_credit", ["OptionalType", ["DataType", "Double"]]], - ["cr_warehouse_sk", ["OptionalType", ["DataType", "Int32"]]] - ] - ] - }, - "catalog_sales": { - "ClusterType": "s3", - "path": "ds/1/catalog_sales/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["cs_bill_addr_sk",["OptionalType",["DataType","Int32"]]], - ["cs_bill_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["cs_bill_customer_sk",["OptionalType",["DataType","Int32"]]], - ["cs_bill_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["cs_call_center_sk",["OptionalType",["DataType","Int32"]]], - ["cs_catalog_page_sk",["OptionalType",["DataType","Int32"]]], - ["cs_coupon_amt",["OptionalType",["DataType","Double"]]], - ["cs_ext_discount_amt",["OptionalType",["DataType","Double"]]], - ["cs_ext_list_price",["OptionalType",["DataType","Double"]]], - ["cs_ext_sales_price",["OptionalType",["DataType","Double"]]], - ["cs_ext_ship_cost",["OptionalType",["DataType","Double"]]], - ["cs_ext_tax",["OptionalType",["DataType","Double"]]], - ["cs_ext_wholesale_cost",["OptionalType",["DataType","Double"]]], - ["cs_item_sk",["OptionalType",["DataType","Int32"]]], - ["cs_list_price",["OptionalType",["DataType","Double"]]], - ["cs_net_paid",["OptionalType",["DataType","Double"]]], - ["cs_net_paid_inc_ship",["OptionalType",["DataType","Double"]]], - ["cs_net_paid_inc_ship_tax",["OptionalType",["DataType","Double"]]], - ["cs_net_paid_inc_tax",["OptionalType",["DataType","Double"]]], - ["cs_net_profit",["OptionalType",["DataType","Double"]]], - ["cs_order_number",["OptionalType",["DataType","Int32"]]], - ["cs_promo_sk",["OptionalType",["DataType","Int32"]]], - ["cs_quantity",["OptionalType",["DataType","Int32"]]], - ["cs_sales_price",["OptionalType",["DataType","Double"]]], - ["cs_ship_addr_sk",["OptionalType",["DataType","Int32"]]], - ["cs_ship_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["cs_ship_customer_sk",["OptionalType",["DataType","Int32"]]], - ["cs_ship_date_sk",["OptionalType",["DataType","Int32"]]], - ["cs_ship_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["cs_ship_mode_sk",["OptionalType",["DataType","Int32"]]], - ["cs_sold_date_sk",["OptionalType",["DataType","Int32"]]], - ["cs_sold_time_sk",["OptionalType",["DataType","Int32"]]], - ["cs_warehouse_sk",["OptionalType",["DataType","Int32"]]], - ["cs_wholesale_cost",["OptionalType",["DataType","Double"]]] - ] - ] - }, - "customer": { - "ClusterType": "s3", - "path": "ds/1/customer/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["c_birth_country",["OptionalType",["DataType","String"]]], - ["c_birth_day",["OptionalType",["DataType","Int32"]]], - ["c_birth_month",["OptionalType",["DataType","Int32"]]], - ["c_birth_year",["OptionalType",["DataType","Int32"]]], - ["c_current_addr_sk",["OptionalType",["DataType","Int32"]]], - ["c_current_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["c_current_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["c_customer_id",["OptionalType",["DataType","String"]]], - ["c_customer_sk",["OptionalType",["DataType","Int32"]]], - ["c_email_address",["OptionalType",["DataType","String"]]], - ["c_first_name",["OptionalType",["DataType","String"]]], - ["c_first_sales_date_sk",["OptionalType",["DataType","Int32"]]], - ["c_first_shipto_date_sk",["OptionalType",["DataType","Int32"]]], - ["c_last_name",["OptionalType",["DataType","String"]]], - ["c_last_review_date",["OptionalType",["DataType","String"]]], - ["c_login",["OptionalType",["DataType","String"]]], - ["c_preferred_cust_flag",["OptionalType",["DataType","String"]]], - ["c_salutation",["OptionalType",["DataType","String"]]] - ] - ] - }, - "customer_address": { - "ClusterType": "s3", - "path": "ds/1/customer_address/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ca_address_id",["OptionalType",["DataType","String"]]], - ["ca_address_sk",["OptionalType",["DataType","Int32"]]], - ["ca_city",["OptionalType",["DataType","String"]]], - ["ca_country",["OptionalType",["DataType","String"]]], - ["ca_county",["OptionalType",["DataType","String"]]], - ["ca_gmt_offset",["OptionalType",["DataType","Double"]]], - ["ca_location_type",["OptionalType",["DataType","String"]]], - ["ca_state",["OptionalType",["DataType","String"]]], - ["ca_street_name",["OptionalType",["DataType","String"]]], - ["ca_street_number",["OptionalType",["DataType","String"]]], - ["ca_street_type",["OptionalType",["DataType","String"]]], - ["ca_suite_number",["OptionalType",["DataType","String"]]], - ["ca_zip",["OptionalType",["DataType","String"]]] - ] - ] - }, - "customer_demographics": { - "ClusterType": "s3", - "path": "ds/1/customer_demographics/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["cd_credit_rating",["OptionalType",["DataType","String"]]], - ["cd_demo_sk",["OptionalType",["DataType","Int32"]]], - ["cd_dep_college_count",["OptionalType",["DataType","Int32"]]], - ["cd_dep_count",["OptionalType",["DataType","Int32"]]], - ["cd_dep_employed_count",["OptionalType",["DataType","Int32"]]], - ["cd_education_status",["OptionalType",["DataType","String"]]], - ["cd_gender",["OptionalType",["DataType","String"]]], - ["cd_marital_status",["OptionalType",["DataType","String"]]], - ["cd_purchase_estimate",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "date_dim": { - "ClusterType": "s3", - "path": "ds/1/date_dim/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["d_current_day",["OptionalType",["DataType","String"]]], - ["d_current_month",["OptionalType",["DataType","String"]]], - ["d_current_quarter",["OptionalType",["DataType","String"]]], - ["d_current_week",["OptionalType",["DataType","String"]]], - ["d_current_year",["OptionalType",["DataType","String"]]], - ["d_date",["OptionalType",["DataType","String"]]], - ["d_date_id",["OptionalType",["DataType","String"]]], - ["d_date_sk",["OptionalType",["DataType","Int32"]]], - ["d_day_name",["OptionalType",["DataType","String"]]], - ["d_dom",["OptionalType",["DataType","Int32"]]], - ["d_dow",["OptionalType",["DataType","Int32"]]], - ["d_first_dom",["OptionalType",["DataType","Int32"]]], - ["d_following_holiday",["OptionalType",["DataType","String"]]], - ["d_fy_quarter_seq",["OptionalType",["DataType","Int32"]]], - ["d_fy_week_seq",["OptionalType",["DataType","Int32"]]], - ["d_fy_year",["OptionalType",["DataType","Int32"]]], - ["d_holiday",["OptionalType",["DataType","String"]]], - ["d_last_dom",["OptionalType",["DataType","Int32"]]], - ["d_month_seq",["OptionalType",["DataType","Int32"]]], - ["d_moy",["OptionalType",["DataType","Int32"]]], - ["d_qoy",["OptionalType",["DataType","Int32"]]], - ["d_quarter_name",["OptionalType",["DataType","String"]]], - ["d_quarter_seq",["OptionalType",["DataType","Int32"]]], - ["d_same_day_lq",["OptionalType",["DataType","Int32"]]], - ["d_same_day_ly",["OptionalType",["DataType","Int32"]]], - ["d_week_seq",["OptionalType",["DataType","Int32"]]], - ["d_weekend",["OptionalType",["DataType","String"]]], - ["d_year",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "household_demographics": { - "ClusterType": "s3", - "path": "ds/1/household_demographics/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["hd_buy_potential",["OptionalType",["DataType","String"]]], - ["hd_demo_sk",["OptionalType",["DataType","Int32"]]], - ["hd_dep_count",["OptionalType",["DataType","Int32"]]], - ["hd_income_band_sk",["OptionalType",["DataType","Int32"]]], - ["hd_vehicle_count",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "income_band": { - "ClusterType": "s3", - "path": "ds/1/income_band/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ib_income_band_sk",["OptionalType",["DataType","Int32"]]], - ["ib_lower_bound",["OptionalType",["DataType","Int32"]]], - ["ib_upper_bound",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "inventory": { - "ClusterType": "s3", - "path": "ds/1/inventory/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["inv_date_sk",["OptionalType",["DataType","Int32"]]], - ["inv_item_sk",["OptionalType",["DataType","Int32"]]], - ["inv_quantity_on_hand",["OptionalType",["DataType","Int32"]]], - ["inv_warehouse_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "item": { - "ClusterType": "s3", - "path": "ds/1/item/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["i_brand",["OptionalType",["DataType","String"]]], - ["i_brand_id",["OptionalType",["DataType","Int32"]]], - ["i_category",["OptionalType",["DataType","String"]]], - ["i_category_id",["OptionalType",["DataType","Int32"]]], - ["i_class",["OptionalType",["DataType","String"]]], - ["i_class_id",["OptionalType",["DataType","Int32"]]], - ["i_color",["OptionalType",["DataType","String"]]], - ["i_container",["OptionalType",["DataType","String"]]], - ["i_current_price",["OptionalType",["DataType","Double"]]], - ["i_formulation",["OptionalType",["DataType","String"]]], - ["i_item_desc",["OptionalType",["DataType","String"]]], - ["i_item_id",["OptionalType",["DataType","String"]]], - ["i_item_sk",["OptionalType",["DataType","Int32"]]], - ["i_manager_id",["OptionalType",["DataType","Int32"]]], - ["i_manufact",["OptionalType",["DataType","String"]]], - ["i_manufact_id",["OptionalType",["DataType","Int32"]]], - ["i_product_name",["OptionalType",["DataType","String"]]], - ["i_rec_end_date",["OptionalType",["DataType","String"]]], - ["i_rec_start_date",["OptionalType",["DataType","String"]]], - ["i_size",["OptionalType",["DataType","String"]]], - ["i_units",["OptionalType",["DataType","String"]]], - ["i_wholesale_cost",["OptionalType",["DataType","Double"]]] - ] - ] - }, - "promotion": { - "ClusterType": "s3", - "path": "ds/1/promotion/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["p_channel_catalog",["OptionalType",["DataType","String"]]], - ["p_channel_demo",["OptionalType",["DataType","String"]]], - ["p_channel_details",["OptionalType",["DataType","String"]]], - ["p_channel_dmail",["OptionalType",["DataType","String"]]], - ["p_channel_email",["OptionalType",["DataType","String"]]], - ["p_channel_event",["OptionalType",["DataType","String"]]], - ["p_channel_press",["OptionalType",["DataType","String"]]], - ["p_channel_radio",["OptionalType",["DataType","String"]]], - ["p_channel_tv",["OptionalType",["DataType","String"]]], - ["p_cost",["OptionalType",["DataType","Double"]]], - ["p_discount_active",["OptionalType",["DataType","String"]]], - ["p_end_date_sk",["OptionalType",["DataType","Int32"]]], - ["p_item_sk",["OptionalType",["DataType","Int32"]]], - ["p_promo_id",["OptionalType",["DataType","String"]]], - ["p_promo_name",["OptionalType",["DataType","String"]]], - ["p_promo_sk",["OptionalType",["DataType","Int32"]]], - ["p_purpose",["OptionalType",["DataType","String"]]], - ["p_response_target",["OptionalType",["DataType","Int32"]]], - ["p_start_date_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "reason": { - "ClusterType": "s3", - "path": "ds/1/reason/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["r_reason_desc",["OptionalType",["DataType","String"]]], - ["r_reason_id",["OptionalType",["DataType","String"]]], - ["r_reason_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "ship_mode": { - "ClusterType": "s3", - "path": "ds/1/ship_mode/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["sm_carrier",["OptionalType",["DataType","String"]]], - ["sm_code",["OptionalType",["DataType","String"]]], - ["sm_contract",["OptionalType",["DataType","String"]]], - ["sm_ship_mode_id",["OptionalType",["DataType","String"]]], - ["sm_ship_mode_sk",["OptionalType",["DataType","Int32"]]], - ["sm_type",["OptionalType",["DataType","String"]]] - ] - ] - }, - "store": { - "ClusterType": "s3", - "path": "ds/1/store/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["s_city",["OptionalType",["DataType","String"]]], - ["s_closed_date_sk",["OptionalType",["DataType","Int32"]]], - ["s_company_id",["OptionalType",["DataType","Int32"]]], - ["s_company_name",["OptionalType",["DataType","String"]]], - ["s_country",["OptionalType",["DataType","String"]]], - ["s_county",["OptionalType",["DataType","String"]]], - ["s_division_id",["OptionalType",["DataType","Int32"]]], - ["s_division_name",["OptionalType",["DataType","String"]]], - ["s_floor_space",["OptionalType",["DataType","Int32"]]], - ["s_geography_class",["OptionalType",["DataType","String"]]], - ["s_gmt_offset",["OptionalType",["DataType","Double"]]], - ["s_hours",["OptionalType",["DataType","String"]]], - ["s_manager",["OptionalType",["DataType","String"]]], - ["s_market_desc",["OptionalType",["DataType","String"]]], - ["s_market_id",["OptionalType",["DataType","Int32"]]], - ["s_market_manager",["OptionalType",["DataType","String"]]], - ["s_number_employees",["OptionalType",["DataType","Int32"]]], - ["s_rec_end_date",["OptionalType",["DataType","String"]]], - ["s_rec_start_date",["OptionalType",["DataType","String"]]], - ["s_state",["OptionalType",["DataType","String"]]], - ["s_store_id",["OptionalType",["DataType","String"]]], - ["s_store_name",["OptionalType",["DataType","String"]]], - ["s_store_sk",["OptionalType",["DataType","Int32"]]], - ["s_street_name",["OptionalType",["DataType","String"]]], - ["s_street_number",["OptionalType",["DataType","String"]]], - ["s_street_type",["OptionalType",["DataType","String"]]], - ["s_suite_number",["OptionalType",["DataType","String"]]], - ["s_tax_precentage",["OptionalType",["DataType","Double"]]], - ["s_zip",["OptionalType",["DataType","String"]]] - ] - ] - }, - "store_returns": { - "ClusterType": "s3", - "path": "ds/1/store_returns/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["sr_addr_sk",["OptionalType",["DataType","Int32"]]], - ["sr_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["sr_customer_sk",["OptionalType",["DataType","Int32"]]], - ["sr_fee",["OptionalType",["DataType","Double"]]], - ["sr_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["sr_item_sk",["OptionalType",["DataType","Int32"]]], - ["sr_net_loss",["OptionalType",["DataType","Double"]]], - ["sr_reason_sk",["OptionalType",["DataType","Int32"]]], - ["sr_refunded_cash",["OptionalType",["DataType","Double"]]], - ["sr_return_amt",["OptionalType",["DataType","Double"]]], - ["sr_return_amt_inc_tax",["OptionalType",["DataType","Double"]]], - ["sr_return_quantity",["OptionalType",["DataType","Int32"]]], - ["sr_return_ship_cost",["OptionalType",["DataType","Double"]]], - ["sr_return_tax",["OptionalType",["DataType","Double"]]], - ["sr_return_time_sk",["OptionalType",["DataType","Int32"]]], - ["sr_returned_date_sk",["OptionalType",["DataType","Int32"]]], - ["sr_reversed_charge",["OptionalType",["DataType","Double"]]], - ["sr_store_credit",["OptionalType",["DataType","Double"]]], - ["sr_store_sk",["OptionalType",["DataType","Int32"]]], - ["sr_ticket_number",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "store_sales": { - "ClusterType": "s3", - "path": "ds/1/store_sales/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ss_addr_sk",["OptionalType",["DataType","Int32"]]], - ["ss_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ss_coupon_amt",["OptionalType",["DataType","Double"]]], - ["ss_customer_sk",["OptionalType",["DataType","Int32"]]], - ["ss_ext_discount_amt",["OptionalType",["DataType","Double"]]], - ["ss_ext_list_price",["OptionalType",["DataType","Double"]]], - ["ss_ext_sales_price",["OptionalType",["DataType","Double"]]], - ["ss_ext_tax",["OptionalType",["DataType","Double"]]], - ["ss_ext_wholesale_cost",["OptionalType",["DataType","Double"]]], - ["ss_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ss_item_sk",["OptionalType",["DataType","Int32"]]], - ["ss_list_price",["OptionalType",["DataType","Double"]]], - ["ss_net_paid",["OptionalType",["DataType","Double"]]], - ["ss_net_paid_inc_tax",["OptionalType",["DataType","Double"]]], - ["ss_net_profit",["OptionalType",["DataType","Double"]]], - ["ss_promo_sk",["OptionalType",["DataType","Int32"]]], - ["ss_quantity",["OptionalType",["DataType","Int32"]]], - ["ss_sales_price",["OptionalType",["DataType","Double"]]], - ["ss_sold_date_sk",["OptionalType",["DataType","Int32"]]], - ["ss_sold_time_sk",["OptionalType",["DataType","Int32"]]], - ["ss_store_sk",["OptionalType",["DataType","Int32"]]], - ["ss_ticket_number",["OptionalType",["DataType","Int32"]]], - ["ss_wholesale_cost",["OptionalType",["DataType","Double"]]] - ] - ] - }, - "time_dim": { - "ClusterType": "s3", - "path": "ds/1/time_dim/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["t_am_pm",["OptionalType",["DataType","String"]]], - ["t_hour",["OptionalType",["DataType","Int32"]]], - ["t_meal_time",["OptionalType",["DataType","String"]]], - ["t_minute",["OptionalType",["DataType","Int32"]]], - ["t_second",["OptionalType",["DataType","Int32"]]], - ["t_shift",["OptionalType",["DataType","String"]]], - ["t_sub_shift",["OptionalType",["DataType","String"]]], - ["t_time",["OptionalType",["DataType","Int32"]]], - ["t_time_id",["OptionalType",["DataType","String"]]], - ["t_time_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "warehouse": { - "ClusterType": "s3", - "path": "ds/1/warehouse/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["w_city",["OptionalType",["DataType","String"]]], - ["w_country",["OptionalType",["DataType","String"]]], - ["w_county",["OptionalType",["DataType","String"]]], - ["w_gmt_offset",["OptionalType",["DataType","Double"]]], - ["w_state",["OptionalType",["DataType","String"]]], - ["w_street_name",["OptionalType",["DataType","String"]]], - ["w_street_number",["OptionalType",["DataType","String"]]], - ["w_street_type",["OptionalType",["DataType","String"]]], - ["w_suite_number",["OptionalType",["DataType","String"]]], - ["w_warehouse_id",["OptionalType",["DataType","String"]]], - ["w_warehouse_name",["OptionalType",["DataType","String"]]], - ["w_warehouse_sk",["OptionalType",["DataType","Int32"]]], - ["w_warehouse_sq_ft",["OptionalType",["DataType","Int32"]]], - ["w_zip",["OptionalType",["DataType","String"]]] - ] - ] - }, - "web_page": { - "ClusterType": "s3", - "path": "ds/1/web_page/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["wp_access_date_sk",["OptionalType",["DataType","Int32"]]], - ["wp_autogen_flag",["OptionalType",["DataType","String"]]], - ["wp_char_count",["OptionalType",["DataType","Int32"]]], - ["wp_creation_date_sk",["OptionalType",["DataType","Int32"]]], - ["wp_customer_sk",["OptionalType",["DataType","Int32"]]], - ["wp_image_count",["OptionalType",["DataType","Int32"]]], - ["wp_link_count",["OptionalType",["DataType","Int32"]]], - ["wp_max_ad_count",["OptionalType",["DataType","Int32"]]], - ["wp_rec_end_date",["OptionalType",["DataType","String"]]], - ["wp_rec_start_date",["OptionalType",["DataType","String"]]], - ["wp_type",["OptionalType",["DataType","String"]]], - ["wp_url",["OptionalType",["DataType","String"]]], - ["wp_web_page_id",["OptionalType",["DataType","String"]]], - ["wp_web_page_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "web_returns": { - "ClusterType": "s3", - "path": "ds/1/web_returns/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["wr_account_credit",["OptionalType",["DataType","Double"]]], - ["wr_fee",["OptionalType",["DataType","Double"]]], - ["wr_item_sk",["OptionalType",["DataType","Int32"]]], - ["wr_net_loss",["OptionalType",["DataType","Double"]]], - ["wr_order_number",["OptionalType",["DataType","Int32"]]], - ["wr_reason_sk",["OptionalType",["DataType","Int32"]]], - ["wr_refunded_addr_sk",["OptionalType",["DataType","Int32"]]], - ["wr_refunded_cash",["OptionalType",["DataType","Double"]]], - ["wr_refunded_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["wr_refunded_customer_sk",["OptionalType",["DataType","Int32"]]], - ["wr_refunded_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["wr_return_amt",["OptionalType",["DataType","Double"]]], - ["wr_return_amt_inc_tax",["OptionalType",["DataType","Double"]]], - ["wr_return_quantity",["OptionalType",["DataType","Int32"]]], - ["wr_return_ship_cost",["OptionalType",["DataType","Double"]]], - ["wr_return_tax",["OptionalType",["DataType","Double"]]], - ["wr_returned_date_sk",["OptionalType",["DataType","Int32"]]], - ["wr_returned_time_sk",["OptionalType",["DataType","Int32"]]], - ["wr_returning_addr_sk",["OptionalType",["DataType","Int32"]]], - ["wr_returning_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["wr_returning_customer_sk",["OptionalType",["DataType","Int32"]]], - ["wr_returning_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["wr_reversed_charge",["OptionalType",["DataType","Double"]]], - ["wr_web_page_sk",["OptionalType",["DataType","Int32"]]] - ] - ] - }, - "web_sales": { - "ClusterType": "s3", - "path": "ds/1/web_sales/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ws_bill_addr_sk",["OptionalType",["DataType","Int32"]]], - ["ws_bill_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ws_bill_customer_sk",["OptionalType",["DataType","Int32"]]], - ["ws_bill_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ws_coupon_amt",["OptionalType",["DataType","Double"]]], - ["ws_ext_discount_amt",["OptionalType",["DataType","Double"]]], - ["ws_ext_list_price",["OptionalType",["DataType","Double"]]], - ["ws_ext_sales_price",["OptionalType",["DataType","Double"]]], - ["ws_ext_ship_cost",["OptionalType",["DataType","Double"]]], - ["ws_ext_tax",["OptionalType",["DataType","Double"]]], - ["ws_ext_wholesale_cost",["OptionalType",["DataType","Double"]]], - ["ws_item_sk",["OptionalType",["DataType","Int32"]]], - ["ws_list_price",["OptionalType",["DataType","Double"]]], - ["ws_net_paid",["OptionalType",["DataType","Double"]]], - ["ws_net_paid_inc_ship",["OptionalType",["DataType","Double"]]], - ["ws_net_paid_inc_ship_tax",["OptionalType",["DataType","Double"]]], - ["ws_net_paid_inc_tax",["OptionalType",["DataType","Double"]]], - ["ws_net_profit",["OptionalType",["DataType","Double"]]], - ["ws_order_number",["OptionalType",["DataType","Int32"]]], - ["ws_promo_sk",["OptionalType",["DataType","Int32"]]], - ["ws_quantity",["OptionalType",["DataType","Int32"]]], - ["ws_sales_price",["OptionalType",["DataType","Double"]]], - ["ws_ship_addr_sk",["OptionalType",["DataType","Int32"]]], - ["ws_ship_cdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ws_ship_customer_sk",["OptionalType",["DataType","Int32"]]], - ["ws_ship_date_sk",["OptionalType",["DataType","Int32"]]], - ["ws_ship_hdemo_sk",["OptionalType",["DataType","Int32"]]], - ["ws_ship_mode_sk",["OptionalType",["DataType","Int32"]]], - ["ws_sold_date_sk",["OptionalType",["DataType","Int32"]]], - ["ws_sold_time_sk",["OptionalType",["DataType","Int32"]]], - ["ws_warehouse_sk",["OptionalType",["DataType","Int32"]]], - ["ws_web_page_sk",["OptionalType",["DataType","Int32"]]], - ["ws_web_site_sk",["OptionalType",["DataType","Int32"]]], - ["ws_wholesale_cost",["OptionalType",["DataType","Double"]]] - ] - ] - }, - "web_site": { - "ClusterType": "s3", - "path": "ds/1/web_site/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["web_city",["OptionalType",["DataType","String"]]], - ["web_class",["OptionalType",["DataType","String"]]], - ["web_close_date_sk",["OptionalType",["DataType","Int32"]]], - ["web_company_id",["OptionalType",["DataType","Int32"]]], - ["web_company_name",["OptionalType",["DataType","String"]]], - ["web_country",["OptionalType",["DataType","String"]]], - ["web_county",["OptionalType",["DataType","String"]]], - ["web_gmt_offset",["OptionalType",["DataType","Double"]]], - ["web_manager",["OptionalType",["DataType","String"]]], - ["web_market_manager",["OptionalType",["DataType","String"]]], - ["web_mkt_class",["OptionalType",["DataType","String"]]], - ["web_mkt_desc",["OptionalType",["DataType","String"]]], - ["web_mkt_id",["OptionalType",["DataType","Int32"]]], - ["web_name",["OptionalType",["DataType","String"]]], - ["web_open_date_sk",["OptionalType",["DataType","Int32"]]], - ["web_rec_end_date",["OptionalType",["DataType","String"]]], - ["web_rec_start_date",["OptionalType",["DataType","String"]]], - ["web_site_id",["OptionalType",["DataType","String"]]], - ["web_site_sk",["OptionalType",["DataType","Int32"]]], - ["web_state",["OptionalType",["DataType","String"]]], - ["web_street_name",["OptionalType",["DataType","String"]]], - ["web_street_number",["OptionalType",["DataType","String"]]], - ["web_street_type",["OptionalType",["DataType","String"]]], - ["web_suite_number",["OptionalType",["DataType","String"]]], - ["web_tax_percentage",["OptionalType",["DataType","Double"]]], - ["web_zip",["OptionalType",["DataType","String"]]] - ] - ] - } -} - diff --git a/ydb/library/yql/tools/dqrun/examples/bindings_tpch.json b/ydb/library/yql/tools/dqrun/examples/bindings_tpch.json deleted file mode 100644 index 64b3d9caff1..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/bindings_tpch.json +++ /dev/null @@ -1,144 +0,0 @@ -{ - "customer": { - "ClusterType": "s3", - "path": "h/1/customer/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["c_acctbal", ["DataType", "Double"]], - ["c_address", ["DataType", "String"]], - ["c_comment", ["DataType", "String"]], - ["c_custkey", ["DataType", "Int32"]], - ["c_mktsegment", ["DataType", "String"]], - ["c_name", ["DataType", "String"]], - ["c_nationkey", ["DataType", "Int32"]], - ["c_phone", ["DataType", "String"]] - ] - ] - }, - "lineitem": { - "ClusterType": "s3", - "path": "h/1/lineitem/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["l_comment", ["DataType", "String"]], - ["l_commitdate", ["DataType", "Date"]], - ["l_discount", ["DataType", "Double"]], - ["l_extendedprice", ["DataType", "Double"]], - ["l_linenumber", ["DataType", "Int32"]], - ["l_linestatus", ["DataType", "String"]], - ["l_orderkey", ["DataType", "Int32"]], - ["l_partkey", ["DataType", "Int32"]], - ["l_quantity", ["DataType", "Double"]], - ["l_receiptdate", ["DataType", "Date"]], - ["l_returnflag", ["DataType", "String"]], - ["l_shipdate", ["DataType", "Date"]], - ["l_shipinstruct", ["DataType", "String"]], - ["l_shipmode", ["DataType", "String"]], - ["l_suppkey", ["DataType", "Int32"]], - ["l_tax", ["DataType", "Double"]] - ] - ] - }, - "nation": { - "ClusterType": "s3", - "path": "h/1/nation/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["n_comment", ["DataType", "String"]], - ["n_name", ["DataType", "String"]], - ["n_nationkey", ["DataType", "Int32"]], - ["n_regionkey", ["DataType", "Int32"]] - ] - ] - }, - "orders": { - "ClusterType": "s3", - "path": "h/1/orders/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["o_clerk", ["DataType", "String"]], - ["o_comment", ["DataType", "String"]], - ["o_custkey", ["DataType", "Int32"]], - ["o_orderdate", ["DataType", "Date"]], - ["o_orderkey", ["DataType", "Int32"]], - ["o_orderpriority", ["DataType", "String"]], - ["o_orderstatus", ["DataType", "String"]], - ["o_shippriority", ["DataType", "Int32"]], - ["o_totalprice", ["DataType", "Double"]] - ] - ] - }, - "part": { - "ClusterType": "s3", - "path": "h/1/part/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["p_brand", ["DataType", "String"]], - ["p_comment", ["DataType", "String"]], - ["p_container", ["DataType", "String"]], - ["p_mfgr", ["DataType", "String"]], - ["p_name", ["DataType", "String"]], - ["p_partkey", ["DataType", "Int32"]], - ["p_retailprice", ["DataType", "Double"]], - ["p_size", ["DataType", "Int32"]], - ["p_type", ["DataType", "String"]] - ] - ] - }, - "partsupp": { - "ClusterType": "s3", - "path": "h/1/partsupp/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ps_availqty", ["DataType", "Int32"]], - ["ps_comment", ["DataType", "String"]], - ["ps_partkey", ["DataType", "Int32"]], - ["ps_suppkey", ["DataType", "Int32"]], - ["ps_supplycost", ["DataType", "Double"]] - ] - ] - }, - "region": { - "ClusterType": "s3", - "path": "h/1/region/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["r_comment", ["DataType", "String"]], - ["r_name", ["DataType", "String"]], - ["r_regionkey", ["DataType", "Int32"]] - ] - ] - }, - "supplier": { - "ClusterType": "s3", - "path": "h/1/supplier/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["s_acctbal", ["DataType", "Double"]], - ["s_address", ["DataType", "String"]], - ["s_comment", ["DataType", "String"]], - ["s_name", ["DataType", "String"]], - ["s_nationkey", ["DataType", "Int32"]], - ["s_phone", ["DataType", "String"]], - ["s_suppkey", ["DataType", "Int32"]] - ] - ] - } -} - diff --git a/ydb/library/yql/tools/dqrun/examples/bindings_tpch_pg.json b/ydb/library/yql/tools/dqrun/examples/bindings_tpch_pg.json deleted file mode 100644 index 28e6d7d0609..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/bindings_tpch_pg.json +++ /dev/null @@ -1,144 +0,0 @@ -{ - "customer": { - "ClusterType": "s3", - "path": "h/1/customer/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["c_acctbal", ["PgType", "numeric"]], - ["c_address", ["PgType", "text"]], - ["c_comment", ["PgType", "text"]], - ["c_custkey", ["PgType", "int4"]], - ["c_mktsegment", ["PgType", "text"]], - ["c_name", ["PgType", "text"]], - ["c_nationkey", ["PgType", "int4"]], - ["c_phone", ["PgType", "text"]] - ] - ] - }, - "lineitem": { - "ClusterType": "s3", - "path": "h/1/lineitem/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["l_comment", ["PgType", "text"]], - ["l_commitdate", ["PgType", "date"]], - ["l_discount", ["PgType", "numeric"]], - ["l_extendedprice", ["PgType", "numeric"]], - ["l_linenumber", ["PgType", "int4"]], - ["l_linestatus", ["PgType", "text"]], - ["l_orderkey", ["PgType", "int4"]], - ["l_partkey", ["PgType", "int4"]], - ["l_quantity", ["PgType", "numeric"]], - ["l_receiptdate", ["PgType", "date"]], - ["l_returnflag", ["PgType", "text"]], - ["l_shipdate", ["PgType", "date"]], - ["l_shipinstruct", ["PgType", "text"]], - ["l_shipmode", ["PgType", "text"]], - ["l_suppkey", ["PgType", "int4"]], - ["l_tax", ["PgType", "numeric"]] - ] - ] - }, - "nation": { - "ClusterType": "s3", - "path": "h/1/nation/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["n_comment", ["PgType", "text"]], - ["n_name", ["PgType", "text"]], - ["n_nationkey", ["PgType", "int4"]], - ["n_regionkey", ["PgType", "int4"]] - ] - ] - }, - "orders": { - "ClusterType": "s3", - "path": "h/1/orders/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["o_clerk", ["PgType", "text"]], - ["o_comment", ["PgType", "text"]], - ["o_custkey", ["PgType", "int4"]], - ["o_orderdate", ["PgType", "date"]], - ["o_orderkey", ["PgType", "int4"]], - ["o_orderpriority", ["PgType", "text"]], - ["o_orderstatus", ["PgType", "text"]], - ["o_shippriority", ["PgType", "int4"]], - ["o_totalprice", ["PgType", "numeric"]] - ] - ] - }, - "part": { - "ClusterType": "s3", - "path": "h/1/part/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["p_brand", ["PgType", "text"]], - ["p_comment", ["PgType", "text"]], - ["p_container", ["PgType", "text"]], - ["p_mfgr", ["PgType", "text"]], - ["p_name", ["PgType", "text"]], - ["p_partkey", ["PgType", "int4"]], - ["p_retailprice", ["PgType", "numeric"]], - ["p_size", ["PgType", "int4"]], - ["p_type", ["PgType", "text"]] - ] - ] - }, - "partsupp": { - "ClusterType": "s3", - "path": "h/1/partsupp/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["ps_availqty", ["PgType", "int4"]], - ["ps_comment", ["PgType", "text"]], - ["ps_partkey", ["PgType", "int4"]], - ["ps_suppkey", ["PgType", "int4"]], - ["ps_supplycost", ["PgType", "numeric"]] - ] - ] - }, - "region": { - "ClusterType": "s3", - "path": "h/1/region/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["r_comment", ["PgType", "text"]], - ["r_name", ["PgType", "text"]], - ["r_regionkey", ["PgType", "int4"]] - ] - ] - }, - "supplier": { - "ClusterType": "s3", - "path": "h/1/supplier/", - "cluster": "yq-tpc-local", - "format": "parquet", - "schema": [ - "StructType", [ - ["s_acctbal", ["PgType", "numeric"]], - ["s_address", ["PgType", "text"]], - ["s_comment", ["PgType", "text"]], - ["s_name", ["PgType", "text"]], - ["s_nationkey", ["PgType", "int4"]], - ["s_phone", ["PgType", "text"]], - ["s_suppkey", ["PgType", "int4"]] - ] - ] - } -} - diff --git a/ydb/library/yql/tools/dqrun/examples/fq.conf b/ydb/library/yql/tools/dqrun/examples/fq.conf deleted file mode 100644 index 6dee04c46e0..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/fq.conf +++ /dev/null @@ -1,21 +0,0 @@ -RowDispatcher { - Enabled: true - TimeoutBeforeStartSessionSec: 2 - MaxSessionUsedMemory: 0 - SendStatusPeriodSec: 10 - WithoutConsumer: true - Coordinator { - CoordinationNodePath: "not_used" - Database { - Endpoint: "not_used:2135" - Database: "/not/used/database" - UseLocalMetadataService: true - UseSsl: true - ClientTimeoutSec: 70 - OperationTimeoutSec: 60 - CancelAfterSec: 60 - } - LocalMode: true - } -} - diff --git a/ydb/library/yql/tools/dqrun/examples/fs.conf b/ydb/library/yql/tools/dqrun/examples/fs.conf deleted file mode 100644 index e3bfbe5a5e3..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/fs.conf +++ /dev/null @@ -1,6 +0,0 @@ -# Use temp directory -#Path: "" -MaxFiles: 1000 -MaxSizeMb: 512 -Threads: 2 -RetryCount: 3 diff --git a/ydb/library/yql/tools/dqrun/examples/gateways.conf b/ydb/library/yql/tools/dqrun/examples/gateways.conf deleted file mode 100644 index c3287d99ce0..00000000000 --- a/ydb/library/yql/tools/dqrun/examples/gateways.conf +++ /dev/null @@ -1,157 +0,0 @@ -Dq { - DefaultSettings { - Name: "EnableComputeActor" - Value: "1" - } - - DefaultSettings { - Name: "ComputeActorType" - Value: "async" - } - - DefaultSettings { - Name: "AnalyzeQuery" - Value: "true" - } - - DefaultSettings { - Name: "MaxTasksPerStage" - Value: "200" - } - - DefaultSettings { - Name: "MaxTasksPerOperation" - Value: "200" - } - - DefaultSettings { - Name: "EnableInsert" - Value: "true" - } - - DefaultSettings { - Name: "_EnablePrecompute" - Value: "true" - } - - DefaultSettings { - Name: "UseAggPhases" - Value: "true" - } - - DefaultSettings { - Name: "HashJoinMode" - Value: "grace" - } - - DefaultSettings { - Name: "UseFastPickleTransport" - Value: "true" - } - - DefaultSettings { - Name: "UseOOBTransport" - Value: "true" - } - - DefaultSettings { - Name: "UseWideChannels" - Value: "true" - } - - DefaultSettings { - Name: "_SkipRevisionCheck" - Value: "true" - } - - DefaultSettings { - Name: "EnableDqReplicate" - Value: "true" - } - DefaultSettings { - Name: "_TableTimeout" - Value: "600000" - } -} - -Generic { - Connector { - Endpoint { - host: "connector.yqv2-dev.cloud.yandex.net" - port: 50051 - } - UseSsl: true - } - - ClusterMapping { - Kind: YDB, - Name: "ydb_dev" - DatabaseId: "etnejle6hb72cdr6aqps" - ServiceAccountId: "my_sa" - ServiceAccountIdSignature: "my_sa_secret_value" - UseSsl: true - Protocol: NATIVE - } - - DefaultSettings { - Name: "DateTimeFormat" - Value: "string" - } -} - -DbResolver { - YdbMvpEndpoint: "https://ydbc.ydb.cloud.yandex.net:8789/ydbc/cloud-prod" -} - -S3 { - ClusterMapping { - Name: "yq-clickbench-local" - Url: "file://./clickbench/" - } - ClusterMapping { - Name: "yq-tpc-local" - Url: "file://./tpc/" - } -} - -HttpGateway { - ConnectionTimeoutSeconds: 15 - RequestTimeoutSeconds: 150 - MaxRetries: 2 - LowSpeedBytesLimit: 16384 - LowSpeedTimeSeconds: 10 - DownloadBufferBytesLimit: 131072 -} - -YqlCore { - Flags { - Name: "_EnableStreamLookupJoin" - } - Flags { - Name: "_EnableMatchRecognize" - } -} - -SqlCore { - ExtendedTranslationFlags: { - Name: "FlexibleTypes" - } - ExtendedTranslationFlags: { - Name: "DisableAnsiOptionalAs" - } - ExtendedTranslationFlags: { - Name: "EmitAggApply" - } -} - -Pq { - ClusterMapping { - Name: "pq" - Endpoint: "localhost:2135" - Database: "local" - ClusterType: CT_DATA_STREAMS - UseSsl: True - SharedReading:True - ReadGroup: "read_group" - } -} diff --git a/ydb/library/yql/tools/dqrun/lib/dqrun_lib.cpp b/ydb/library/yql/tools/dqrun/lib/dqrun_lib.cpp deleted file mode 100644 index 9c7cb4d7de5..00000000000 --- a/ydb/library/yql/tools/dqrun/lib/dqrun_lib.cpp +++ /dev/null @@ -1,560 +0,0 @@ -#include "dqrun_lib.h" - -#include <yt/yql/providers/yt/gateway/file/yql_yt_file.h> -#include <yt/yql/providers/yt/gateway/native/yql_yt_native.h> -#include <yt/yql/providers/yt/provider/yql_yt_provider.h> -#include <yt/yql/providers/yt/gateway/file/yql_yt_file_comp_nodes.h> -#include <yt/yql/providers/yt/comp_nodes/dq/dq_yt_factory.h> -#include <yt/yql/providers/yt/mkql_dq/yql_yt_dq_transform.h> - -#include <yql/essentials/providers/common/provider/yql_provider_names.h> -#include <yql/essentials/providers/common/comp_nodes/yql_factory.h> -#include <yql/essentials/core/dq_integration/transform/yql_dq_task_transform.h> -#include <yql/essentials/minikql/comp_nodes/mkql_factories.h> -#include <yql/essentials/parser/pg_wrapper/interface/comp_factory.h> - -#include <ydb/core/fq/libs/db_id_async_resolver_impl/database_resolver.h> -#include <ydb/core/fq/libs/db_id_async_resolver_impl/db_async_resolver_impl.h> -#include <ydb/core/fq/libs/db_id_async_resolver_impl/mdb_endpoint_generator.h> -#include <ydb/core/fq/libs/shared_resources/interface/shared_resources.h> -#include <ydb/core/fq/libs/config/protos/fq_config.pb.h> -#include <ydb/core/fq/libs/init/init.h> -#include <ydb/library/yql/providers/clickhouse/provider/yql_clickhouse_provider.h> -#include <ydb/library/yql/providers/clickhouse/actors/yql_ch_source_factory.h> -#include <ydb/library/yql/providers/ydb/provider/yql_ydb_provider.h> -#include <ydb/library/yql/providers/ydb/comp_nodes/yql_ydb_factory.h> -#include <ydb/library/yql/providers/ydb/comp_nodes/yql_ydb_dq_transform.h> -#include <ydb/library/yql/providers/ydb/actors/yql_ydb_source_factory.h> -#include <ydb/library/yql/providers/s3/provider/yql_s3_provider.h> -#include <ydb/library/yql/providers/s3/actors/yql_s3_actors_factory_impl.h> -#include <ydb/library/yql/providers/pq/provider/yql_pq_provider.h> -#include <ydb/library/yql/providers/pq/gateway/dummy/yql_pq_dummy_gateway.h> -#include <ydb/library/yql/providers/pq/gateway/native/yql_pq_gateway.h> -#include <ydb/library/yql/providers/pq/async_io/dq_pq_read_actor.h> -#include <ydb/library/yql/providers/pq/async_io/dq_pq_write_actor.h> -#include <ydb/library/yql/providers/solomon/provider/yql_solomon_provider.h> -#include <ydb/library/yql/providers/solomon/gateway/yql_solomon_gateway.h> -#include <ydb/library/yql/providers/solomon/actors/dq_solomon_read_actor.h> -#include <ydb/library/yql/providers/generic/provider/yql_generic_provider.h> -#include <ydb/library/yql/providers/generic/actors/yql_generic_provider_factories.h> -#include <ydb/library/yql/providers/dq/interface/yql_dq_task_preprocessor.h> -#include <ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h> -#include <ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.h> -#include <ydb/library/yql/providers/dq/provider/yql_dq_provider.h> -#include <ydb/library/yql/providers/dq/helper/yql_dq_helper_impl.h> -#include <ydb/library/yql/providers/yt/dq_task_preprocessor/yql_yt_dq_task_preprocessor.h> -#include <ydb/library/yql/providers/yt/actors/yql_yt_provider_factories.h> -#include <ydb/library/yql/providers/common/http_gateway/yql_http_default_retry_policy.h> -#include <ydb/library/yql/dq/comp_nodes/yql_common_dq_factory.h> -#include <ydb/library/yql/dq/opt/dq_opt_join_cbo_factory.h> -#include <ydb/library/yql/dq/transform/yql_common_dq_transform.h> -#include <ydb/library/yql/dq/actors/input_transforms/dq_input_transform_lookup_factory.h> -#include <ydb/library/yql/utils/bindings/utils.h> -#include <ydb/library/actors/http/http_proxy.h> - -#include <util/string/builder.h> -#include <util/stream/output.h> -#include <util/stream/file.h> - -#ifdef PROFILE_MEMORY_ALLOCATIONS -#include <library/cpp/lfalloc/alloc_profiler/profiler.h> -#endif - -#ifdef __unix__ -#include <sys/resource.h> -#endif - -namespace NYql { - -const std::initializer_list<TString> SUPPORTED_GATEWAYS = { - TString{DqProviderName}, - TString{YtProviderName}, - TString{SolomonProviderName}, - TString{ClickHouseProviderName}, - TString{YdbProviderName}, - TString{PqProviderName}, - TString{S3ProviderName}, - TString{GenericProviderName}, - TString{"hybrid"}, -}; - -TDqRunTool::TDqRunTool() - : TYtRunTool("dqrun") -{ - GetRunOptions().SetSupportedGateways(SUPPORTED_GATEWAYS); - GetRunOptions().GatewayTypes.insert(SUPPORTED_GATEWAYS); - - GetRunOptions().EnableCredentials = true; - GetRunOptions().EnableQPlayer = true; - GetRunOptions().ResultStream = &Cout; - GetRunOptions().ResultsFormat = NYson::EYsonFormat::Text; - GetRunOptions().CustomTests = true; - GetRunOptions().Verbosity = TLOG_INFO; - - GetRunOptions().AddOptExtension([this](NLastGetopt::TOpts& opts) { - - opts.AddLongOption("dq-host", "Dq host") - .RequiredArgument("HOST") - .Handler1T<TString>([this](const TString& host) { - DqHost_ = host; - }); - opts.AddLongOption("dq-port", "Dq port") - .RequiredArgument("PORT") - .Handler1T<int>([this](int port) { - DqPort_ = port; - }); - opts.AddLongOption("analyze-query", "enable analyze query") - .Optional() - .NoArgument() - .SetFlag(&AnalyzeQuery_); - opts.AddLongOption("no-force-dq", "don't set force dq mode") - .Optional() - .NoArgument() - .SetFlag(&NoForceDq_); - opts.AddLongOption("enable-spilling", "Enable disk spilling") - .NoArgument() - .SetFlag(&EnableSpilling_); - opts.AddLongOption("tmp-dir", "directory for temporary tables") - .StoreResult<TString>(&TmpDir_); - opts.AddLongOption('t', "table", "Table mapping").RequiredArgument("table@file") - .KVHandler([&](TString name, TString path) { - if (name.empty() || path.empty()) { - throw yexception() << "Incorrect table mapping, expected form table@file, e.g. [email protected]"; - } - TablesMapping_[name] = path; - }, '@'); - opts.AddLongOption('C', "cluster", "Cluster to service mapping").RequiredArgument("name@service") - .KVHandler([&](TString cluster, TString provider) { - if (cluster.empty() || provider.empty()) { - throw yexception() << "Incorrect service mapping, expected form cluster@provider, e.g. plato@yt"; - } - AddClusterMapping(std::move(cluster), std::move(provider)); - }, '@'); - opts.AddLongOption("dq-threads", "DQ gateway threads") - .Optional() - .RequiredArgument("COUNT") - .StoreResult(&DqThreads_); - opts.AddLongOption("fq-cfg", "federated query configuration file") - .Optional() - .RequiredArgument("FILE") - .Handler1T<TString>([this](const TString& file) { - FqConfig_ = TFacadeRunOptions::ParseProtoConfig<NFq::NConfig::TConfig>(file); - }); - opts.AddLongOption("metrics", "Print execution metrics") - .Optional() - .OptionalArgument("FILE") - .Handler1T<TString>([this](const TString& file) { - if (file) { - MetricsStreamHolder_ = MakeHolder<TFileOutput>(file); - MetricsStream_ = MetricsStreamHolder_.Get(); - } else { - MetricsStream_ = &Cerr; - } - }); - opts.AddLongOption('E', "emulate-yt", "Emulate YT tables") - .Optional() - .NoArgument() - .SetFlag(&EmulateYt_); - opts.AddLongOption("emulate-pq", "Emulate YDS with local file") - .RequiredArgument("topic@file") - .KVHandler([&](TString name, TString path) { - if (name.empty() || path.empty()) { - throw yexception() << "Incorrect topic mapping, expected form topic@file"; - } - TopicMapping_[name] = path; - }, '@'); - - opts.AddLongOption("bindings-file", "Bindings File") - .Optional() - .RequiredArgument("FILE") - .Handler1T<TString>([this](const TString& file) { - TFileInput input(file); - LoadBindings(GetRunOptions().Bindings, input.ReadAll()); - }); - opts.AddLongOption("token-accessor-endpoint", "Network address of Token Accessor service in format grpc(s)://host:port") - .Optional() - .RequiredArgument("ENDPOINT") - .Handler1T<TString>([this](const TString& endpoint) { - TVector<TString> ss = StringSplitter(endpoint).SplitByString("://"); - if (ss.size() != 2) { - throw yexception() << "Invalid tokenAccessorEndpoint: " << endpoint; - } - CredentialsFactory_ = NYql::CreateSecuredServiceAccountCredentialsOverTokenAccessorFactory(ss[1], ss[0] == "grpcs", ""); - }); - }); - - GetRunOptions().AddOptHandler([this](const NLastGetopt::TOptsParseResult& res) { - Y_UNUSED(res); - - if (EmulateYt_ && DqPort_) { - throw yexception() << "Remote DQ instance cannot work with the emulated YT cluster"; - } - if (EmulateYt_) { - GetRunOptions().GatewayTypes.emplace(YtProviderName); - } - - GetRunOptions().EnableResultPosition = !EmulateYt_; - GetRunOptions().UseRepeatableRandomAndTimeProviders = EmulateYt_; - - if (!GetRunOptions().GatewaysConfig) { - GetRunOptions().GatewaysConfig = MakeHolder<TGatewaysConfig>(); - } - - auto dqConfig = GetRunOptions().GatewaysConfig->MutableDq(); - - if (AnalyzeQuery_) { - auto* setting = dqConfig->AddDefaultSettings(); - setting->SetName("AnalyzeQuery"); - setting->SetValue("1"); - } - - if (EnableSpilling_) { - auto* setting = dqConfig->AddDefaultSettings(); - setting->SetName("SpillingEngine"); - setting->SetValue("file"); - } - if (EmulateYt_) { - AddClusterMapping(TString{"plato"}, TString{YtProviderName}); - // Clusters from gateways are filled in YtRunLib - } - - if (GetRunOptions().GatewaysConfig->HasClickHouse()) { - FillClusterMapping(GetRunOptions().GatewaysConfig->GetClickHouse(), TString{ClickHouseProviderName}); - } - - if (GetRunOptions().GatewaysConfig->HasGeneric()) { - FillClusterMapping(GetRunOptions().GatewaysConfig->GetGeneric(), TString{GenericProviderName}); - } - - if (GetRunOptions().GatewaysConfig->HasYdb()) { - FillClusterMapping(GetRunOptions().GatewaysConfig->GetYdb(), TString{YdbProviderName}); - } - - if (GetRunOptions().GatewaysConfig->HasS3()) { - GetRunOptions().GatewaysConfig->MutableS3()->SetAllowLocalFiles(true); - FillClusterMapping(GetRunOptions().GatewaysConfig->GetS3(), TString{S3ProviderName}); - } - - if (GetRunOptions().GatewaysConfig->HasPq()) { - FillClusterMapping(GetRunOptions().GatewaysConfig->GetPq(), TString{PqProviderName}); - } - - if (GetRunOptions().GatewaysConfig->HasSolomon()) { - FillClusterMapping(GetRunOptions().GatewaysConfig->GetSolomon(), TString{SolomonProviderName}); - } - - if (!FqConfig_) { - FqConfig_ = MakeHolder<NFq::NConfig::TConfig>(); - } - - GetRunOptions().SqlFlags["DqEngineEnable"] = {}; - if (!AnalyzeQuery_ && !NoForceDq_) { - GetRunOptions().SqlFlags["DqEngineForce"] = {}; - } - - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasClickHouse()) { - return GetClickHouseDataProviderInitializer(GetHttpGateway()); - } - return {}; - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasGeneric()) { - return GetGenericDataProviderInitializer(GetGenericClient(), GetDbResolver(), CredentialsFactory_); - } - return {}; - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasYdb()) { - return GetYdbDataProviderInitializer(GetYdbDriver()); - } - return {}; - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasS3()) { - return GetS3DataProviderInitializer(GetHttpGateway(), nullptr); - } - return {}; - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasPq() || !TopicMapping_.empty()) { - return GetPqDataProviderInitializer(GetPqGateway(), false, GetDbResolver()); - } - return {}; - }); - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - if (GetRunOptions().GatewaysConfig->HasSolomon()) { - return GetSolomonDataProviderInitializer(CreateSolomonGateway(GetRunOptions().GatewaysConfig->GetSolomon()), nullptr, false); - } - return {}; - }); - - - AddProviderFactory([this]() -> NYql::TDataProviderInitializer { - auto compFactory = CreateCompNodeFactory(); - TIntrusivePtr<IDqGateway> dqGateway; - if (DqPort_) { - dqGateway = CreateDqGateway(DqHost_.GetOrElse("localhost"), *DqPort_); - } else { - std::function<NActors::IActor*(void)> metricsPusherFactory = {}; - dqGateway = CreateLocalDqGateway(GetFuncRegistry().Get(), compFactory, CreateDqTaskTransformFactory(), CreateDqTaskPreprocessorFactories(), - EnableSpilling_, CreateAsyncIoFactory(), DqThreads_, GetMetricsRegistry(), metricsPusherFactory); - } - - return GetDqDataProviderInitializer(&CreateDqExecTransformer, dqGateway, compFactory, {}, GetFileStorage()); - }); -} - -TDqRunTool::~TDqRunTool() { - try { - if (Driver_) { - Driver_->Stop(true); - } - } catch (...) { - Cerr << "Error while stopping YDB driver: " << CurrentExceptionMessage() << Endl; - } -} - -IYtGateway::TPtr TDqRunTool::CreateYtGateway() { - if (EmulateYt_) { - return CreateYtFileGateway(GetYtFileServices()); - } - return TYtRunTool::CreateYtGateway(); -} - -IOptimizerFactory::TPtr TDqRunTool::CreateCboFactory() { - return NDq::MakeCBOOptimizerFactory(); -} - -IDqHelper::TPtr TDqRunTool::CreateDqHelper() { - return MakeDqHelper(); -} - -NYdb::TDriver TDqRunTool::GetYdbDriver() { - if (!Driver_) { - const auto driverConfig = NYdb::TDriverConfig().SetLog(std::unique_ptr<TLogBackend>(CreateLogBackend("cerr").Release())); - Driver_.ConstructInPlace(driverConfig); - } - return *Driver_; -} - -NFile::TYtFileServices::TPtr TDqRunTool::GetYtFileServices() { - if (!YtFileServices_) { - YtFileServices_ = NFile::TYtFileServices::Make(GetFuncRegistry().Get(), TablesMapping_, GetFileStorage(), TmpDir_, KeepTemp_); - } - return YtFileServices_; -} - -IHTTPGateway::TPtr TDqRunTool::GetHttpGateway() { - if (!HttpGateway_) { - HttpGateway_ = IHTTPGateway::Make(GetRunOptions().GatewaysConfig->HasHttpGateway() ? &GetRunOptions().GatewaysConfig->GetHttpGateway() : nullptr); - } - return HttpGateway_; -} - -IMetricsRegistryPtr TDqRunTool::GetMetricsRegistry() { - if (!MetricsRegistry_) { - MetricsRegistry_ = CreateMetricsRegistry(GetSensorsGroupFor(NSensorComponent::kDq)); - } - return MetricsRegistry_; -} - -NActors::NLog::EPriority TDqRunTool::YqlToActorsLogLevel(NYql::NLog::ELevel yqlLevel) { - switch (yqlLevel) { - case NYql::NLog::ELevel::FATAL: - return NActors::NLog::PRI_CRIT; - case NYql::NLog::ELevel::ERROR: - return NActors::NLog::PRI_ERROR; - case NYql::NLog::ELevel::WARN: - return NActors::NLog::PRI_WARN; - case NYql::NLog::ELevel::INFO: - return NActors::NLog::PRI_INFO; - case NYql::NLog::ELevel::DEBUG: - return NActors::NLog::PRI_DEBUG; - case NYql::NLog::ELevel::TRACE: - return NActors::NLog::PRI_TRACE; - default: - ythrow yexception() << "unexpected level: " << int(yqlLevel); - } -} - -struct TDqRunTool::TActorIds { - NActors::TActorId DatabaseResolver; - NActors::TActorId HttpProxy; -}; - -std::tuple<std::unique_ptr<TActorSystemManager>, TDqRunTool::TActorIds> TDqRunTool::RunActorSystem() { - auto actorSystemManager = std::make_unique<TActorSystemManager>(GetMetricsRegistry(), YqlToActorsLogLevel(NYql::NLog::ELevelHelpers::FromInt(GetRunOptions().Verbosity))); - TActorIds actorIds; - - // Run actor system only if necessary - auto needActorSystem = GetRunOptions().GatewaysConfig->HasGeneric() || GetRunOptions().GatewaysConfig->HasDbResolver(); - if (!needActorSystem) { - return std::make_tuple(std::move(actorSystemManager), std::move(actorIds)); - } - - // One can modify actor system setup via actorSystemManager->ApplySetupModifier(). - // TODO: https://st.yandex-team.ru/YQL-16131 - // This will be useful for DQ Gateway initialization refactoring. - actorSystemManager->Start(); - - // Actor system is initialized; start actor registration. - if (needActorSystem) { - auto httpProxy = NHttp::CreateHttpProxy(); - actorIds.HttpProxy = actorSystemManager->GetActorSystem()->Register(httpProxy); - - auto databaseResolver = NFq::CreateDatabaseResolver(actorIds.HttpProxy, CredentialsFactory_); - actorIds.DatabaseResolver = actorSystemManager->GetActorSystem()->Register(databaseResolver); - } - - return std::make_tuple(std::move(actorSystemManager), std::move(actorIds)); -} - -NYql::IDatabaseAsyncResolver::TPtr TDqRunTool::GetDbResolver() { - if (!DbResolver_ && GetRunOptions().GatewaysConfig->HasDbResolver()) { - std::unique_ptr<TActorSystemManager> actorSystemManager; - TActorIds actorIds; - std::tie(actorSystemManager, actorIds) = RunActorSystem(); - - DbResolver_ = std::make_shared<NFq::TDatabaseAsyncResolverImpl>( - actorSystemManager->GetActorSystem(), - actorIds.DatabaseResolver, - GetRunOptions().GatewaysConfig->GetDbResolver().GetYdbMvpEndpoint(), - GetRunOptions().GatewaysConfig->HasGeneric() ? GetRunOptions().GatewaysConfig->GetGeneric().GetMdbGateway() : "", - NFq::MakeMdbEndpointGeneratorGeneric(false) - ); - } - return DbResolver_; -} - -NConnector::IClient::TPtr TDqRunTool::GetGenericClient() { - if (!GenericClient_ && GetRunOptions().GatewaysConfig->HasGeneric()) { - GenericClient_ = NConnector::MakeClientGRPC(GetRunOptions().GatewaysConfig->GetGeneric()); - } - return GenericClient_; -} - -IPqGateway::TPtr TDqRunTool::GetPqGateway() { - if (!PqGateway_ && (GetRunOptions().GatewaysConfig->HasPq() || !TopicMapping_.empty())) { - if (!TopicMapping_.empty()) { - auto fileGateway = MakeIntrusive<TDummyPqGateway>(); - for (const auto& pair: TopicMapping_) { - fileGateway->AddDummyTopic(TDummyTopic("pq", pair.first, pair.second)); - } - PqGateway_ = fileGateway; - } else { - TPqGatewayServices pqServices( - GetYdbDriver(), - nullptr, - nullptr, // credentials factory - std::make_shared<TPqGatewayConfig>(GetRunOptions().GatewaysConfig->GetPq()), - GetFuncRegistry().Get() - ); - PqGateway_ = CreatePqNativeGateway(pqServices); - } - } - return PqGateway_; -} - -NKikimr::NMiniKQL::TComputationNodeFactory TDqRunTool::CreateCompNodeFactory() { - TVector<NKikimr::NMiniKQL::TComputationNodeFactory> factories = { - GetCommonDqFactory(), - NKikimr::NMiniKQL::GetYqlFactory(), - GetPgFactory() - }; - - if (GetRunOptions().GatewaysConfig->HasYt() || EmulateYt_) { - factories.push_back(GetDqYtFactory()); - } - if (EmulateYt_) { - factories.push_back(GetYtFileFactory(GetYtFileServices())); - } - if (GetRunOptions().GatewaysConfig->HasYdb()) { - factories.push_back(GetDqYdbFactory(GetYdbDriver())); - } - return NKikimr::NMiniKQL::GetCompositeWithBuiltinFactory(factories); -} - -TTaskTransformFactory TDqRunTool::CreateDqTaskTransformFactory() { - TVector<TTaskTransformFactory> factories = { - CreateCommonDqTaskTransformFactory() - }; - - if (GetRunOptions().GatewaysConfig->HasYt()) { - factories.push_back(CreateYtDqTaskTransformFactory()); - } - if (GetRunOptions().GatewaysConfig->HasYdb()) { - factories.push_back(CreateYdbDqTaskTransformFactory()); - } - - return CreateCompositeTaskTransformFactory(std::move(factories)); -} - -TDqTaskPreprocessorFactoryCollection TDqRunTool::CreateDqTaskPreprocessorFactories() { - TDqTaskPreprocessorFactoryCollection factories; - - if (GetRunOptions().GatewaysConfig->HasYt() || EmulateYt_) { - factories.push_back(NDq::CreateYtDqTaskPreprocessorFactory(EmulateYt_, GetFuncRegistry())); - } - return factories; -} - -NDq::IDqAsyncIoFactory::TPtr TDqRunTool::CreateAsyncIoFactory() { - auto factory = MakeIntrusive<NYql::NDq::TDqAsyncIoFactory>(); - RegisterDqInputTransformLookupActorFactory(*factory); - if (EmulateYt_) { - RegisterYtLookupActorFactory(*factory, GetYtFileServices(), *GetFuncRegistry()); - } - if (GetRunOptions().GatewaysConfig->HasPq() || !TopicMapping_.empty()) { - RegisterDqPqReadActorFactory(*factory, GetYdbDriver(), nullptr, GetPqGateway()); - RegisterDqPqWriteActorFactory(*factory, GetYdbDriver(), nullptr, GetPqGateway()); - } - if (GetRunOptions().GatewaysConfig->HasYdb()) { - RegisterYdbReadActorFactory(*factory, GetYdbDriver(), nullptr); - } - if (GetRunOptions().GatewaysConfig->HasSolomon()) { - RegisterDQSolomonReadActorFactory(*factory, nullptr); - } - if (GetRunOptions().GatewaysConfig->HasClickHouse()) { - RegisterClickHouseReadActorFactory(*factory, nullptr, GetHttpGateway()); - } - if (GetRunOptions().GatewaysConfig->HasGeneric()) { - RegisterGenericProviderFactories(*factory, CredentialsFactory_, GetGenericClient()); - } - if (GetRunOptions().GatewaysConfig->HasS3()) { - auto s3ActorsFactory = NYql::NDq::CreateS3ActorsFactory(); - - const size_t requestTimeout = GetRunOptions().GatewaysConfig->HasHttpGateway() && GetRunOptions().GatewaysConfig->GetHttpGateway().HasRequestTimeoutSeconds() - ? GetRunOptions().GatewaysConfig->GetHttpGateway().GetRequestTimeoutSeconds() - : 100; - const size_t maxRetries = GetRunOptions().GatewaysConfig->HasHttpGateway() && GetRunOptions().GatewaysConfig->GetHttpGateway().HasMaxRetries() - ? GetRunOptions().GatewaysConfig->GetHttpGateway().GetMaxRetries() - : 2; - - s3ActorsFactory->RegisterS3WriteActorFactory(*factory, nullptr, GetHttpGateway(), GetHTTPDefaultRetryPolicy()); - s3ActorsFactory->RegisterS3ReadActorFactory(*factory, nullptr, GetHttpGateway(), GetHTTPDefaultRetryPolicy(TDuration::Seconds(requestTimeout), maxRetries)); - } - - return factory; -} - - -TProgram::TStatus TDqRunTool::DoRunProgram(TProgramPtr program) { -#ifdef PROFILE_MEMORY_ALLOCATIONS - NAllocProfiler::StartAllocationSampling(true); -#endif - const TProgram::TStatus status = TYtRunTool::DoRunProgram(program); -#ifdef PROFILE_MEMORY_ALLOCATIONS - NAllocProfiler::StopAllocationSampling(Cout); -#endif - return status; -} - -} // NYql diff --git a/ydb/library/yql/tools/dqrun/lib/dqrun_lib.h b/ydb/library/yql/tools/dqrun/lib/dqrun_lib.h deleted file mode 100644 index 6bd29b6d744..00000000000 --- a/ydb/library/yql/tools/dqrun/lib/dqrun_lib.h +++ /dev/null @@ -1,94 +0,0 @@ -#pragma once - -#include <yt/yql/tools/ytrun/lib/ytrun_lib.h> -#include <yt/yql/providers/yt/gateway/file/yql_yt_file_services.h> - -#include <yql/essentials/tools/yql_facade_run/yql_facade_run.h> -#include <yql/essentials/providers/common/metrics/metrics_registry.h> -#include <yql/essentials/core/cbo/cbo_optimizer_new.h> -#include <yql/essentials/core/dq_integration/yql_dq_helper.h> -#include <yql/essentials/core/dq_integration/transform/yql_dq_task_transform.h> -#include <yql/essentials/sql/settings/translation_settings.h> -#include <yql/essentials/minikql/computation/mkql_computation_node.h> -#include <yql/essentials/utils/log/log_level.h> - -#include <ydb/library/actors/core/log_iface.h> -#include <ydb/library/yql/dq/actors/compute/dq_compute_actor_async_io.h> -#include <ydb/library/yql/providers/common/db_id_async_resolver/db_async_resolver.h> -#include <ydb/library/yql/providers/common/http_gateway/yql_http_gateway.h> -#include <ydb/library/yql/providers/common/token_accessor/client/factory.h> -#include <ydb/library/yql/providers/generic/connector/libcpp/client.h> -#include <ydb/library/yql/providers/dq/interface/yql_dq_task_preprocessor.h> -#include <ydb/library/yql/providers/pq/gateway/abstract/yql_pq_gateway.h> -#include <ydb/library/yql/utils/actor_system/manager.h> -#include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/driver/driver.h> - -#include <util/generic/string.h> -#include <util/generic/hash.h> -#include <util/generic/maybe.h> - -namespace NFq::NConfig { - class TConfig; -} - -namespace NYql { - -class TDqRunTool: public TYtRunTool { -public: - TDqRunTool(); - ~TDqRunTool(); - -protected: - struct TActorIds; - - TProgram::TStatus DoRunProgram(TProgramPtr program) override; - IOptimizerFactory::TPtr CreateCboFactory() override; - IDqHelper::TPtr CreateDqHelper() override; - IYtGateway::TPtr CreateYtGateway() override; - - NYdb::TDriver GetYdbDriver(); - NYql::NFile::TYtFileServices::TPtr GetYtFileServices(); - IHTTPGateway::TPtr GetHttpGateway(); - IMetricsRegistryPtr GetMetricsRegistry(); - static NActors::NLog::EPriority YqlToActorsLogLevel(NYql::NLog::ELevel yqlLevel); - std::tuple<std::unique_ptr<TActorSystemManager>, TActorIds> RunActorSystem(); - NYql::IDatabaseAsyncResolver::TPtr GetDbResolver(); - NConnector::IClient::TPtr GetGenericClient(); - IPqGateway::TPtr GetPqGateway(); - NKikimr::NMiniKQL::TComputationNodeFactory CreateCompNodeFactory(); - NYql::TTaskTransformFactory CreateDqTaskTransformFactory(); - NYql::TDqTaskPreprocessorFactoryCollection CreateDqTaskPreprocessorFactories(); - NYql::NDq::IDqAsyncIoFactory::TPtr CreateAsyncIoFactory(); - -protected: - bool AnalyzeQuery_ = false; - bool NoForceDq_ = false; - bool EmulateYt_ = false; - TMaybe<TString> DqHost_; - TMaybe<int> DqPort_; - int DqThreads_ = 16; - bool EnableSpilling_ = false; - - IOutputStream* MetricsStream_ = nullptr; - THolder<IOutputStream> MetricsStreamHolder_; - - TString MrJobBin_; - TString MrJobUdfsDir_; - bool KeepTemp_ = false; - TString TmpDir_; - THashMap<TString, TString> TablesMapping_; - - THashMap<TString, TString> TopicMapping_; - THolder<NFq::NConfig::TConfig> FqConfig_; - - ISecuredServiceAccountCredentialsFactory::TPtr CredentialsFactory_; - TMaybe<NYdb::TDriver> Driver_; - NFile::TYtFileServices::TPtr YtFileServices_; - IHTTPGateway::TPtr HttpGateway_; - IMetricsRegistryPtr MetricsRegistry_; - NYql::IDatabaseAsyncResolver::TPtr DbResolver_; - NConnector::IClient::TPtr GenericClient_; - IPqGateway::TPtr PqGateway_; -}; - -} diff --git a/ydb/library/yql/tools/dqrun/lib/lsan.supp b/ydb/library/yql/tools/dqrun/lib/lsan.supp deleted file mode 100644 index e23b02a4276..00000000000 --- a/ydb/library/yql/tools/dqrun/lib/lsan.supp +++ /dev/null @@ -1,6 +0,0 @@ -leak:Py_InitializeEx -leak:RegisterYqlPythonUdf -leak:PyBytes_FromStringAndSize -leak:PyUnicode_New -leak:_PyObject_MallocWithType -leak:init_interp_main diff --git a/ydb/library/yql/tools/dqrun/lib/ya.make b/ydb/library/yql/tools/dqrun/lib/ya.make deleted file mode 100644 index 0afaf5afdc5..00000000000 --- a/ydb/library/yql/tools/dqrun/lib/ya.make +++ /dev/null @@ -1,86 +0,0 @@ -LIBRARY() - -SRCS( - dqrun_lib.cpp -) - -PEERDIR( - yt/yql/tools/ytrun/lib - yt/yql/providers/yt/gateway/file - yt/yql/providers/yt/gateway/native - yt/yql/providers/yt/provider - yt/yql/providers/yt/gateway/file - yt/yql/providers/yt/comp_nodes/dq - yt/yql/providers/yt/mkql_dq - - yql/essentials/providers/common/provider - yql/essentials/providers/common/comp_nodes - yql/essentials/providers/common/metrics - yql/essentials/core/dq_integration - yql/essentials/core/dq_integration/transform - yql/essentials/minikql/computation - yql/essentials/parser/pg_wrapper/interface - yql/essentials/tools/yql_facade_run - yql/essentials/sql/settings - yql/essentials/utils/log - - ydb/core/fq/libs/db_id_async_resolver_impl - ydb/core/fq/libs/shared_resources/interface - ydb/core/fq/libs/config/protos - ydb/core/fq/libs/init - ydb/core/fq/libs/actors - - ydb/library/yql/providers/clickhouse/provider - ydb/library/yql/providers/clickhouse/actors - - ydb/library/yql/providers/ydb/provider - ydb/library/yql/providers/ydb/comp_nodes - ydb/library/yql/providers/ydb/actors - - ydb/library/yql/providers/s3/provider - ydb/library/yql/providers/s3/actors - - ydb/library/yql/providers/pq/provider - ydb/library/yql/providers/pq/gateway/abstract - ydb/library/yql/providers/pq/async_io - - ydb/library/yql/providers/solomon/provider - ydb/library/yql/providers/solomon/gateway - ydb/library/yql/providers/solomon/actors - - ydb/library/yql/providers/generic/connector/libcpp - ydb/library/yql/providers/generic/provider - ydb/library/yql/providers/generic/actors - - ydb/library/yql/providers/dq/interface - ydb/library/yql/providers/dq/local_gateway - ydb/library/yql/providers/dq/provider/exec - ydb/library/yql/providers/dq/provider - ydb/library/yql/providers/dq/helper - - ydb/library/yql/providers/yt/dq_task_preprocessor - ydb/library/yql/providers/yt/actors - - ydb/library/yql/providers/common/token_accessor/client - ydb/library/yql/providers/common/http_gateway - ydb/library/yql/providers/common/db_id_async_resolver - - ydb/library/yql/dq/opt - ydb/library/yql/dq/transform - ydb/library/yql/dq/actors/compute - ydb/library/yql/dq/actors/input_transforms - - ydb/library/yql/utils/bindings - ydb/library/yql/utils/actor_system - - ydb/library/actors/core - ydb/library/actors/http -) - -YQL_LAST_ABI_VERSION() - -SUPPRESSIONS( - lsan.supp -) - -END() diff --git a/ydb/library/yql/tools/dqrun/ya.make b/ydb/library/yql/tools/dqrun/ya.make deleted file mode 100644 index 601c06b94b9..00000000000 --- a/ydb/library/yql/tools/dqrun/ya.make +++ /dev/null @@ -1,51 +0,0 @@ -IF (NOT OS_WINDOWS) - PROGRAM() - - IF (PROFILE_MEMORY_ALLOCATIONS) - ALLOCATOR(LF_DBG) - CFLAGS(-DPROFILE_MEMORY_ALLOCATIONS) - ELSE() - IF (OS_LINUX AND NOT DISABLE_TCMALLOC) - ALLOCATOR(TCMALLOC_256K) - ELSE() - ALLOCATOR(J) - ENDIF() - ENDIF() - - - IF (OOM_HELPER) - PEERDIR(yql/essentials/utils/oom_helper) - ENDIF() - - SRCS( - dqrun.cpp - ) - - PEERDIR( - ydb/library/yql/tools/dqrun/lib - - yt/yql/providers/yt/codec/codegen - yt/yql/providers/yt/comp_nodes/llvm16 - yt/yql/providers/yt/comp_nodes/dq/llvm16 - yql/essentials/minikql/invoke_builtins/llvm16 - yql/essentials/minikql/comp_nodes/llvm16 - yql/essentials/parser/pg_wrapper - yql/essentials/public/udf/service/exception_policy - yql/essentials/sql/pg - - library/cpp/lfalloc/alloc_profiler - - ydb/library/yql/udfs/common/clickhouse/client - ydb/library/yql/dq/comp_nodes/llvm16 - ydb/library/yql/providers/pq/gateway/dummy - ydb/public/sdk/cpp/src/client/persqueue_public/codecs - ) - - YQL_LAST_ABI_VERSION() - - END() -ELSE() - LIBRARY() - - END() -ENDIF() diff --git a/ydb/library/yql/tools/ya.make b/ydb/library/yql/tools/ya.make index 42e4fc5cd90..22cbca53f6d 100644 --- a/ydb/library/yql/tools/ya.make +++ b/ydb/library/yql/tools/ya.make @@ -1,4 +1,3 @@ RECURSE( - dqrun solomon_emulator ) |
