summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorSergey Uzhakov <[email protected]>2026-07-15 14:54:14 +0300
committerGitHub <[email protected]>2026-07-15 14:54:14 +0300
commit1dd74a25bf627151fced1f37a0045205b9e48405 (patch)
tree8a0efe433e6d83a0bb095470bf25b7852fb9f38c
parentcf6219c66d10cb909b70c559cf9de6a115cf3fab (diff)
remove dqrun and unused dq-related code (#46501)
-rw-r--r--ydb/library/yql/providers/dq/actors/ya.make1
-rw-r--r--ydb/library/yql/providers/dq/common/yql_dq_common.h6
-rw-r--r--ydb/library/yql/providers/dq/config/config.proto193
-rw-r--r--ydb/library/yql/providers/dq/config/ya.make8
-rw-r--r--ydb/library/yql/providers/dq/local_gateway/ut/ya.make23
-rw-r--r--ydb/library/yql/providers/dq/local_gateway/ut/yql_dq_gateway_local_ut.cpp56
-rw-r--r--ydb/library/yql/providers/dq/local_gateway/ya.make22
-rw-r--r--ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.cpp302
-rw-r--r--ydb/library/yql/providers/dq/local_gateway/yql_dq_gateway_local.h24
-rw-r--r--ydb/library/yql/providers/dq/provider/exec/ya.make29
-rw-r--r--ydb/library/yql/providers/dq/provider/exec/yql_dq_exectransformer.cpp5
-rw-r--r--ydb/library/yql/providers/dq/provider/ut/ya.make41
-rw-r--r--ydb/library/yql/providers/dq/provider/ya.make35
-rw-r--r--ydb/library/yql/providers/dq/provider/yql_dq_control.cpp197
-rw-r--r--ydb/library/yql/providers/dq/provider/yql_dq_control.h46
-rw-r--r--ydb/library/yql/providers/dq/provider/yql_dq_gateway.cpp802
-rw-r--r--ydb/library/yql/providers/dq/provider/yql_dq_gateway.h7
-rw-r--r--ydb/library/yql/providers/dq/provider/yql_dq_provider_ut.cpp365
-rw-r--r--ydb/library/yql/providers/dq/service/grpc_service.cpp863
-rw-r--r--ydb/library/yql/providers/dq/service/grpc_service.h47
-rw-r--r--ydb/library/yql/providers/dq/service/grpc_session.cpp124
-rw-r--r--ydb/library/yql/providers/dq/service/grpc_session.h59
-rw-r--r--ydb/library/yql/providers/dq/service/interconnect_helpers.cpp317
-rw-r--r--ydb/library/yql/providers/dq/service/interconnect_helpers.h51
-rw-r--r--ydb/library/yql/providers/dq/service/service_node.cpp167
-rw-r--r--ydb/library/yql/providers/dq/service/service_node.h48
-rw-r--r--ydb/library/yql/providers/dq/service/ya.make33
-rw-r--r--ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.cpp31
-rw-r--r--ydb/library/yql/providers/dq/stats_collector/pool_stats_collector.h18
-rw-r--r--ydb/library/yql/providers/dq/stats_collector/ya.make20
-rw-r--r--ydb/library/yql/providers/dq/ya.make2
-rw-r--r--ydb/library/yql/tools/dqrun/.gitignore1
-rw-r--r--ydb/library/yql/tools/dqrun/README.md132
-rw-r--r--ydb/library/yql/tools/dqrun/dqrun.cpp12
-rw-r--r--ydb/library/yql/tools/dqrun/examples/bindings_tpcds.json668
-rw-r--r--ydb/library/yql/tools/dqrun/examples/bindings_tpch.json144
-rw-r--r--ydb/library/yql/tools/dqrun/examples/bindings_tpch_pg.json144
-rw-r--r--ydb/library/yql/tools/dqrun/examples/fq.conf21
-rw-r--r--ydb/library/yql/tools/dqrun/examples/fs.conf6
-rw-r--r--ydb/library/yql/tools/dqrun/examples/gateways.conf157
-rw-r--r--ydb/library/yql/tools/dqrun/lib/dqrun_lib.cpp560
-rw-r--r--ydb/library/yql/tools/dqrun/lib/dqrun_lib.h94
-rw-r--r--ydb/library/yql/tools/dqrun/lib/lsan.supp6
-rw-r--r--ydb/library/yql/tools/dqrun/lib/ya.make86
-rw-r--r--ydb/library/yql/tools/dqrun/ya.make51
-rw-r--r--ydb/library/yql/tools/ya.make1
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(&params);
- 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(&params);
- } 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
)