summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorVitalii Gridnev <[email protected]>2024-03-07 14:12:35 +0300
committerGitHub <[email protected]>2024-03-07 14:12:35 +0300
commit2fd02381daa62921b520bfd65e0975dfcdb1e557 (patch)
tree36f8ba932f0a766769103f6c468a9401cd015762
parent56825aad40a640357b0b1d8e42de10c7bad30955 (diff)
sessions system table (#2460)
-rw-r--r--ydb/core/kqp/common/events/events.h17
-rw-r--r--ydb/core/kqp/common/simple/kqp_event_ids.h4
-rw-r--r--ydb/core/kqp/proxy_service/kqp_proxy_service.cpp110
-rw-r--r--ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h17
-rw-r--r--ydb/core/kqp/ut/sysview/kqp_sys_view_ut.cpp46
-rw-r--r--ydb/core/protos/kqp.proto23
-rw-r--r--ydb/core/sys_view/common/schema.cpp1
-rw-r--r--ydb/core/sys_view/common/schema.h17
-rw-r--r--ydb/core/sys_view/scan.cpp5
-rw-r--r--ydb/core/sys_view/sessions/sessions.cpp261
-rw-r--r--ydb/core/sys_view/sessions/sessions.h15
-rw-r--r--ydb/core/sys_view/sessions/ya.make21
-rw-r--r--ydb/core/sys_view/ut_kqp.cpp4
-rw-r--r--ydb/core/sys_view/ya.make1
14 files changed, 539 insertions, 3 deletions
diff --git a/ydb/core/kqp/common/events/events.h b/ydb/core/kqp/common/events/events.h
index 6a8607c10a9..cda9ae27a2e 100644
--- a/ydb/core/kqp/common/events/events.h
+++ b/ydb/core/kqp/common/events/events.h
@@ -2,6 +2,7 @@
#include "process_response.h"
#include "query.h"
+
#include <ydb/core/kqp/common/simple/kqp_event_ids.h>
#include <ydb/core/protos/kqp.pb.h>
#include <ydb/core/kqp/common/compilation/events.h>
@@ -56,6 +57,22 @@ struct TEvKqp {
using TEvQueryResponse = NPrivateEvents::TEvQueryResponse;
+ struct TEvListSessionsRequest: public TEventPB<TEvListSessionsRequest, NKikimrKqp::TEvListSessionsRequest,
+ TKqpEvents::EvListSessionsRequest>
+ {};
+
+ struct TEvListSessionsResponse: public TEventPB<TEvListSessionsResponse, NKikimrKqp::TEvListSessionsResponse,
+ TKqpEvents::EvListSessionsResponse>
+ {};
+
+ struct TEvListProxyNodesRequest : public TEventLocal<TEvListProxyNodesRequest, TKqpEvents::EvListProxyNodesRequest>
+ {};
+
+ struct TEvListProxyNodesResponse : public TEventLocal<TEvListProxyNodesResponse, TKqpEvents::EvListProxyNodesResponse>
+ {
+ std::vector<ui32> ProxyNodes;
+ };
+
struct TEvCreateSessionResponse : public TEventPB<TEvCreateSessionResponse,
NKikimrKqp::TEvCreateSessionResponse, TKqpEvents::EvCreateSessionResponse> {};
diff --git a/ydb/core/kqp/common/simple/kqp_event_ids.h b/ydb/core/kqp/common/simple/kqp_event_ids.h
index 5adbc8de18a..9b1632e46df 100644
--- a/ydb/core/kqp/common/simple/kqp_event_ids.h
+++ b/ydb/core/kqp/common/simple/kqp_event_ids.h
@@ -41,6 +41,10 @@ struct TKqpEvents {
EvParseRequest,
EvParseResponse,
EvSplitResponse,
+ EvListSessionsRequest,
+ EvListSessionsResponse,
+ EvListProxyNodesRequest,
+ EvListProxyNodesResponse
};
static_assert (EvCompileInvalidateRequest + 1 == EvAbortExecution);
diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp b/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp
index 6871dada0e0..1f07b14756e 100644
--- a/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp
+++ b/ydb/core/kqp/proxy_service/kqp_proxy_service.cpp
@@ -27,6 +27,7 @@
#include <ydb/core/kqp/compute_actor/kqp_compute_actor.h>
#include <ydb/core/mon/mon.h>
#include <ydb/library/ydb_issue/issue_helpers.h>
+#include <ydb/core/sys_view/common/schema.h>
#include <ydb/library/yql/utils/actor_log/log.h>
#include <ydb/library/yql/core/services/mounts/yql_mounts.h>
@@ -671,6 +672,10 @@ public:
return;
}
+ if (sessionInfo) {
+ LocalSessions->AttachQueryText(sessionInfo, ev->Get()->GetQuery());
+ }
+
TActorId targetId;
if (sessionInfo) {
targetId = sessionInfo->WorkerId;
@@ -1315,6 +1320,8 @@ public:
hFunc(NKqp::TEvCancelScriptExecutionOperation, Handle);
hFunc(TEvInterconnect::TEvNodeConnected, Handle);
hFunc(TEvInterconnect::TEvNodeDisconnected, Handle);
+ hFunc(TEvKqp::TEvListSessionsRequest, Handle);
+ hFunc(TEvKqp::TEvListProxyNodesRequest, Handle);
default:
Y_ABORT("TKqpProxyService: unexpected event type: %" PRIx32 " event: %s",
ev->GetTypeRewrite(), ev->ToString().data());
@@ -1626,6 +1633,109 @@ private:
}
}
+ void Handle(TEvKqp::TEvListSessionsRequest::TPtr& ev) {
+ auto result = std::make_unique<TEvKqp::TEvListSessionsResponse>();
+ auto startIt = LocalSessions->GetOrderedLowerBound(ev->Get()->Record.GetSessionIdStart());
+ auto endIt = LocalSessions->GetOrderedEnd();
+ i32 freeSpace = ev->Get()->Record.GetFreeSpace();
+
+ struct TFieldsMap {
+ bool NeedSessionId = false;
+ bool NeedQueryText = false;
+ bool NeedStatus = false;
+ bool NeedNodeId = false;
+
+ explicit TFieldsMap(const ::google::protobuf::RepeatedField<ui32>& columns) {
+ for(const auto& column: columns) {
+ switch(column) {
+ case NKikimr::NSysView::Schema::QuerySessions::SessionId::ColumnId:
+ NeedSessionId = true;
+ break;
+ case NKikimr::NSysView::Schema::QuerySessions::Status::ColumnId:
+ NeedStatus = true;
+ break;
+ case NKikimr::NSysView::Schema::QuerySessions::QueryText::ColumnId:
+ NeedQueryText = true;
+ break;
+ case NKikimr::NSysView::Schema::QuerySessions::NodeId::ColumnId:
+ NeedNodeId = true;
+ break;
+ default: {
+ Y_UNREACHABLE();
+ }
+ }
+ }
+ }
+ };
+
+ TFieldsMap fieldsMap(ev->Get()->Record.GetColumns());
+
+ const TString until = ev->Get()->Record.GetSessionIdEnd();
+ bool finished = false;
+
+ while(startIt != endIt && freeSpace > 0) {
+ auto* sessionInfo = startIt->second;
+
+ if (!until.empty()) {
+ if (sessionInfo->SessionId > until) {
+ finished = true;
+ break;
+ }
+
+ if (!ev->Get()->Record.GetSessionIdEndInclusive() && until == sessionInfo->SessionId) {
+ finished = true;
+ break;
+ }
+ }
+
+ auto* sessionProto = result->Record.AddSessions();
+ sessionProto->SetSessionId(sessionInfo->SessionId);
+ // last executed query or currently running query.
+ if (fieldsMap.NeedQueryText) {
+ sessionProto->SetQueryText(sessionInfo->QueryText);
+ }
+
+ if (fieldsMap.NeedStatus) {
+ if (LocalSessions->IsSessionIdle(sessionInfo)) {
+ sessionProto->SetStatus("IDLE");
+ } else {
+ sessionProto->SetStatus("RUNNING");
+ }
+ }
+
+ freeSpace -= sessionProto->ByteSizeLong();
+ ++startIt;
+ }
+
+ if (startIt == endIt) {
+ finished = true;
+ }
+
+ result->Record.SetNodeId(SelfId().NodeId());
+ if (finished) {
+ result->Record.SetFinished(true);
+ } else {
+ result->Record.SetContinuationToken(startIt->first);
+ result->Record.SetFinished(false);
+ }
+
+ Send(ev->Sender, result.release(), 0, ev->Cookie);
+ }
+
+ void Handle(TEvKqp::TEvListProxyNodesRequest::TPtr& ev) {
+ auto result = std::make_unique<TEvKqp::TEvListProxyNodesResponse>();
+ result->ProxyNodes.reserve(PeerProxyNodeResources.size());
+ for(const auto& resource: PeerProxyNodeResources) {
+ result->ProxyNodes.push_back(resource.GetNodeId());
+ }
+
+ if (result->ProxyNodes.size() < 1) {
+ result->ProxyNodes.push_back(SelfId().NodeId());
+ }
+
+ Send(ev->Sender, result.release(), 0, ev->Cookie);
+ }
+
private:
NKikimrConfig::TLogConfig LogConfig;
NKikimrConfig::TTableServiceConfig TableServiceConfig;
diff --git a/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h b/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h
index be252f3c339..d8bce4e058c 100644
--- a/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h
+++ b/ydb/core/kqp/proxy_service/kqp_proxy_service_impl.h
@@ -88,6 +88,8 @@ struct TKqpSessionInfo {
TNodeId AttachedNodeId;
TActorId AttachedRpcId;
bool PgWire;
+ TString QueryText;
+ bool Ready = true;
TKqpSessionInfo(const TString& sessionId, const TActorId& workerId,
const TString& database, TKqpDbCountersPtr dbCounters, std::vector<i32>&& pos,
@@ -108,6 +110,7 @@ struct TKqpSessionInfo {
class TLocalSessionsRegistry {
THashMap<TString, TKqpSessionInfo> LocalSessions;
+ std::map<TString, TKqpSessionInfo*> OrderedSessions;
THashMap<TActorId, TString> TargetIdIndex;
THashSet<TString> ShutdownInFlightSessions;
THashMap<TString, ui32> SessionsCountPerDatabase;
@@ -130,6 +133,10 @@ public:
return actors.insert(sessionInfo).second;
}
+ void AttachQueryText(const TKqpSessionInfo* sessionInfo, const TString& queryText) {
+ const_cast<TKqpSessionInfo*>(sessionInfo)->QueryText = queryText;
+ }
+
TKqpSessionInfo* Create(const TString& sessionId, const TActorId& workerId,
const TString& database, TKqpDbCountersPtr dbCounters, bool supportsBalancing,
TDuration idleDuration, bool pgWire = false)
@@ -146,6 +153,7 @@ public:
auto result = LocalSessions.emplace(sessionId,
TKqpSessionInfo(sessionId, workerId, database, dbCounters, std::move(pos),
NActors::TActivationContext::Monotonic() + idleDuration, IdleSessions.end(), pgWire));
+ OrderedSessions.emplace(sessionId, &result.first->second);
SessionsCountPerDatabase[database]++;
Y_ABORT_UNLESS(result.second, "Duplicate session id!");
TargetIdIndex.emplace(workerId, sessionId);
@@ -239,6 +247,14 @@ public:
return ShutdownInFlightSessions.size();
}
+ std::map<TString, TKqpSessionInfo*>::const_iterator GetOrderedLowerBound(const TString& continuation) const {
+ return OrderedSessions.lower_bound(continuation);
+ }
+
+ std::map<TString, TKqpSessionInfo*>::const_iterator GetOrderedEnd() const {
+ return OrderedSessions.end();
+ }
+
std::pair<TNodeId, TActorId> Erase(const TString& sessionId) {
auto it = LocalSessions.find(sessionId);
auto result = std::make_pair<TNodeId, TActorId>(0, TActorId());
@@ -268,6 +284,7 @@ public:
}
}
+ OrderedSessions.erase(sessionId);
LocalSessions.erase(it);
}
diff --git a/ydb/core/kqp/ut/sysview/kqp_sys_view_ut.cpp b/ydb/core/kqp/ut/sysview/kqp_sys_view_ut.cpp
index 095f15e77f7..e0b2a23f8ab 100644
--- a/ydb/core/kqp/ut/sysview/kqp_sys_view_ut.cpp
+++ b/ydb/core/kqp/ut/sysview/kqp_sys_view_ut.cpp
@@ -44,6 +44,52 @@ Y_UNIT_TEST_SUITE(KqpSystemView) {
CompareYson(R"([[["::1"];["/Root/KeyValue"];[2u]]])", StreamResultToYson(it));
}
+ Y_UNIT_TEST(Sessions) {
+ TKikimrRunner kikimr;
+ auto client = kikimr.GetQueryClient();
+ const size_t sessionsCount = 50;
+ std::vector<NYdb::NQuery::TSession> sessionsSet;
+ for(ui32 i = 0; i < sessionsCount; i++) {
+ sessionsSet.emplace_back(std::move(client.GetSession().GetValueSync().GetSession()));
+ }
+
+ Cerr << kikimr.GetTestServer().GetRuntime()->GetNodeId() << Endl;
+
+ ui32 nodeId = kikimr.GetTestServer().GetRuntime()->GetNodeId();
+
+ std::sort(sessionsSet.begin(), sessionsSet.end(), [](const NYdb::NQuery::TSession& a, const NYdb::NQuery::TSession& b){
+ return a.GetId() < b.GetId();
+ });
+
+ std::vector<TString> stringParts;
+ for(ui32 i = 0; i < sessionsCount - 1; i++) {
+ stringParts.push_back(Sprintf("[[\"%s\"];[%du];[\"\"]];", sessionsSet[i].GetId().data(), nodeId));
+ }
+
+ TString otherSessions = JoinSeq("\n", stringParts);
+
+ {
+ auto result = sessionsSet.back().ExecuteQuery(R"(--!syntax_v1
+select SessionId, NodeId, QueryText from `/Root/.sys/query_sessions` order by SessionId;)", NYdb::NQuery::TTxControl::NoTx()).GetValueSync();
+
+ UNIT_ASSERT_C(result.IsSuccess(), result.GetIssues().ToString());
+
+ CompareYson(Sprintf(R"([
+ %s
+ [["%s"];[%du];["--!syntax_v1\nselect SessionId, NodeId, QueryText from `/Root/.sys/query_sessions` order by SessionId;"]]
+ ])", otherSessions.data(), sessionsSet.back().GetId().data(), nodeId), FormatResultSetYson(result.GetResultSet(0)));
+ }
+
+ {
+ auto result = sessionsSet.back().ExecuteQuery(Sprintf(R"(--!syntax_v1
+select SessionId, NodeId, QueryText from `/Root/.sys/query_sessions` WHERE StartsWith(SessionId, "ydb://session/3?node_id=%d");)", nodeId), NYdb::NQuery::TTxControl::NoTx()).GetValueSync();
+ CompareYson(Sprintf(R"([
+ %s
+ [["%s"];[%du];["--!syntax_v1\nselect SessionId, NodeId, QueryText from `/Root/.sys/query_sessions` WHERE StartsWith(SessionId, \"ydb://session/3?node_id=%d\");"]]
+ ])", otherSessions.data(), sessionsSet.back().GetId().data(), nodeId, nodeId), FormatResultSetYson(result.GetResultSet(0)));
+ }
+ }
+
Y_UNIT_TEST(PartitionStatsSimple) {
TKikimrRunner kikimr;
auto client = kikimr.GetTableClient();
diff --git a/ydb/core/protos/kqp.proto b/ydb/core/protos/kqp.proto
index e87afd4a75f..b57ff7fb9f0 100644
--- a/ydb/core/protos/kqp.proto
+++ b/ydb/core/protos/kqp.proto
@@ -307,6 +307,29 @@ message TEvProcessResponse {
optional bool WorkerIsClosing = 4 [default = false];
}
+message TSessionInfo {
+ optional string SessionId = 1;
+ optional string Status = 2;
+ optional string QueryText = 3;
+}
+
+message TEvListSessionsRequest {
+ optional string SessionIdStart = 1;
+ optional bool SessionIdStartInclusive = 2;
+ optional string SessionIdEnd = 3;
+ optional bool SessionIdEndInclusive = 4;
+ repeated uint32 Columns = 5;
+ optional int64 FreeSpace = 6;
+ optional int64 Limit = 7;
+}
+
+message TEvListSessionsResponse {
+ optional int64 NodeId = 1;
+ repeated TSessionInfo Sessions = 2;
+ optional bool Finished = 3;
+ optional string ContinuationToken = 4;
+}
+
message TKqpSetting {
required string Name = 1;
required string Value = 2;
diff --git a/ydb/core/sys_view/common/schema.cpp b/ydb/core/sys_view/common/schema.cpp
index 80ef81d4dec..7ff74b75c0a 100644
--- a/ydb/core/sys_view/common/schema.cpp
+++ b/ydb/core/sys_view/common/schema.cpp
@@ -203,6 +203,7 @@ private:
RegisterSystemView<Schema::QueryStats>(TopQueriesByCpuTime1HourName);
RegisterSystemView<Schema::QueryStats>(TopQueriesByRequestUnits1MinuteName);
RegisterSystemView<Schema::QueryStats>(TopQueriesByRequestUnits1HourName);
+ RegisterSystemView<Schema::QuerySessions>(QuerySessions);
RegisterDomainSystemView<Schema::PDisks>(PDisksName);
RegisterDomainSystemView<Schema::VSlots>(VSlotsName);
diff --git a/ydb/core/sys_view/common/schema.h b/ydb/core/sys_view/common/schema.h
index 8eb8fe8fdd9..633a3c712e0 100644
--- a/ydb/core/sys_view/common/schema.h
+++ b/ydb/core/sys_view/common/schema.h
@@ -10,6 +10,7 @@ namespace NSysView {
constexpr TStringBuf PartitionStatsName = "partition_stats";
constexpr TStringBuf NodesName = "nodes";
+constexpr TStringBuf QuerySessions = "query_sessions";
constexpr TStringBuf TopQueriesByDuration1MinuteName = "top_queries_by_duration_one_minute";
constexpr TStringBuf TopQueriesByDuration1HourName = "top_queries_by_duration_one_hour";
@@ -410,7 +411,7 @@ struct Schema : NIceDb::Schema {
Rows,
RawBytes,
PortionId,
- ChunkIdx,
+ ChunkIdx,
EntityName,
InternalEntityId,
BlobId,
@@ -463,6 +464,20 @@ struct Schema : NIceDb::Schema {
IndexSize,
InFlightTxCount>;
};
+
+ struct QuerySessions : Table<13> {
+ struct SessionId : Column<1, NScheme::NTypeIds::Utf8> {};
+ struct NodeId : Column<2, NScheme::NTypeIds::Uint32> {};
+ struct Status : Column<3, NScheme::NTypeIds::Utf8> {};
+ struct QueryText : Column<4, NScheme::NTypeIds::Utf8> {};
+
+ using TKey = TableKey<SessionId>;
+ using TColumns = TableColumns<
+ SessionId,
+ NodeId,
+ Status,
+ QueryText>;
+ };
};
bool MaybeSystemViewPath(const TVector<TString>& path);
diff --git a/ydb/core/sys_view/scan.cpp b/ydb/core/sys_view/scan.cpp
index a8ce212a49a..bf65fa4e10b 100644
--- a/ydb/core/sys_view/scan.cpp
+++ b/ydb/core/sys_view/scan.cpp
@@ -7,6 +7,7 @@
#include <ydb/core/sys_view/nodes/nodes.h>
#include <ydb/core/sys_view/query_stats/query_stats.h>
#include <ydb/core/sys_view/query_stats/query_metrics.h>
+#include <ydb/core/sys_view/sessions/sessions.h>
#include <ydb/core/sys_view/storage/pdisks.h>
#include <ydb/core/sys_view/storage/vslots.h>
#include <ydb/core/sys_view/storage/groups.h>
@@ -161,6 +162,10 @@ THolder<NActors::IActor> CreateSystemViewScan(const NActors::TActorId& ownerId,
return CreateNodesScan(ownerId, scanId, tableId, tableRange, columns);
}
+ if (tableId.SysViewInfo == QuerySessions) {
+ return CreateSessionsScan(ownerId, scanId, tableId, tableRange, columns);
+ }
+
if (tableId.SysViewInfo == TopQueriesByDuration1MinuteName ||
tableId.SysViewInfo == TopQueriesByDuration1HourName ||
tableId.SysViewInfo == TopQueriesByReadBytes1MinuteName ||
diff --git a/ydb/core/sys_view/sessions/sessions.cpp b/ydb/core/sys_view/sessions/sessions.cpp
new file mode 100644
index 00000000000..33f26cadda4
--- /dev/null
+++ b/ydb/core/sys_view/sessions/sessions.cpp
@@ -0,0 +1,261 @@
+#include "sessions.h"
+
+
+#include <ydb/library/actors/core/interconnect.h>
+#include <ydb/core/sys_view/common/events.h>
+#include <ydb/core/sys_view/common/schema.h>
+#include <ydb/core/sys_view/common/scan_actor_base_impl.h>
+#include <ydb/core/node_whiteboard/node_whiteboard.h>
+#include <ydb/core/kqp/common/simple/services.h>
+#include <ydb/core/kqp/common/events/events.h>
+
+#include <ydb/library/yql/dq/actors/compute/dq_compute_actor.h>
+
+#include <ydb/library/actors/core/interconnect.h>
+#include <ydb/library/actors/interconnect/interconnect.h>
+#include <ydb/library/actors/core/hfunc.h>
+
+namespace NKikimr::NSysView {
+
+using namespace NActors;
+using namespace NNodeWhiteboard;
+
+class TSessionsScan : public NKikimr::NSysView::TScanActorBase<TSessionsScan> {
+public:
+ using TBase = NKikimr::NSysView::TScanActorBase<TSessionsScan>;
+
+ static constexpr auto ActorActivityType() {
+ return NKikimrServices::TActivity::KQP_SYSTEM_VIEW_SCAN;
+ }
+
+ using TNodeInfo = NKikimrKqp::TSessionInfo;
+ using TExtractor = std::function<TCell(const TNodeInfo&, ui32)>;
+ using TSchema = Schema::QuerySessions;
+
+ struct TExtractorsMap : public THashMap<NTable::TTag, TExtractor> {
+ TExtractorsMap() {
+ insert({TSchema::NodeId::ColumnId, [] (const TNodeInfo&, ui32 nodeId) {
+ return TCell::Make<ui32>(nodeId);
+ }});
+
+ insert({TSchema::SessionId::ColumnId, [] (const TNodeInfo& info, ui32) {
+ return TCell(info.GetSessionId().data(), info.GetSessionId().size());
+ }});
+ insert({TSchema::QueryText::ColumnId, [] (const TNodeInfo& info, ui32) {
+ return TCell(info.GetQueryText().data(), info.GetQueryText().size());
+ }});
+
+ insert({TSchema::Status::ColumnId, [] (const TNodeInfo& info, ui32) {
+ return TCell(info.GetStatus().data(), info.GetStatus().size());
+ }});
+ }
+ };
+
+ TSessionsScan(const NActors::TActorId& ownerId, ui32 scanId, const TTableId& tableId,
+ const TTableRange& tableRange, const TArrayRef<NMiniKQL::TKqpComputeContextBase::TColumn>& columns)
+ : TBase(ownerId, scanId, tableId, tableRange, columns)
+ {
+ const auto& cellsFrom = TableRange.From.GetCells();
+ if (cellsFrom.size() == 1 && !cellsFrom[0].IsNull()) {
+ SessionIdFrom = cellsFrom[0].AsBuf();
+ SessionIdFromInclusive = TableRange.FromInclusive;
+ }
+
+ const auto& cellsTo = TableRange.To.GetCells();
+ if (cellsTo.size() == 1 && !cellsTo[0].IsNull()) {
+ SessionIdTo = cellsTo[0].AsBuf();
+ SessionIdToInclusive = TableRange.ToInclusive;
+ }
+
+ static TExtractorsMap extractors;
+ ColumnsExtractors.reserve(Columns.size());
+ ColumnsToRead.reserve(Columns.size());
+ for(const auto& column: Columns) {
+ auto it = extractors.find(column.Tag);
+ ColumnsToRead.push_back(column.Tag);
+ if (it != extractors.end()) {
+ ColumnsExtractors.push_back(it->second);
+ } else {
+ MissingSchemaColumns.push_back(column.Tag);
+ }
+ }
+ }
+
+ STFUNC(StateScan) {
+ switch (ev->GetTypeRewrite()) {
+ hFunc(NKqp::TEvKqpCompute::TEvScanDataAck, Handle);
+ hFunc(TEvents::TEvUndelivered, Undelivered);
+ hFunc(NActors::TEvInterconnect::TEvNodeConnected, Connected);
+ hFunc(NActors::TEvInterconnect::TEvNodeDisconnected, Disconnected);
+ hFunc(NKqp::TEvKqp::TEvAbortExecution, HandleAbortExecution);
+ cFunc(TEvents::TEvWakeup::EventType, HandleTimeout);
+ cFunc(TEvents::TEvPoison::EventType, PassAway);
+ hFunc(NKqp::TEvKqp::TEvListSessionsResponse, Handle);
+ hFunc(NKqp::TEvKqp::TEvListProxyNodesResponse, Handle);
+ default:
+ LOG_CRIT(*TlsActivationContext, NKikimrServices::SYSTEM_VIEWS,
+ "NSysView::TSessionsScan: unexpected event 0x%08" PRIx32, ev->GetTypeRewrite());
+ }
+ }
+
+private:
+ void ProceedToScan() override {
+ Become(&TSessionsScan::StateScan);
+ if (!MissingSchemaColumns.empty()) {
+ TStringBuilder message;
+ message << "Missing schema column tags: ";
+ bool first = true;
+ for(ui32 tag: MissingSchemaColumns) {
+ if (!first) {
+ message << ", ";
+ }
+ if (first) first = false;
+ message << tag;
+ }
+
+ ReplyErrorAndDie(Ydb::StatusIds::INTERNAL_ERROR, message);
+ return;
+ }
+
+ if (AckReceived) {
+ StartScan();
+ }
+ }
+
+ void StartScan() {
+ if (IsEmptyRange) {
+ ReplyEmptyAndDie();
+ return;
+ }
+
+ if (!PendingNodesInitialized) {
+ Send(NKqp::MakeKqpProxyID(SelfId().NodeId()), new NKikimr::NKqp::TEvKqp::TEvListProxyNodesRequest());
+ return;
+ }
+
+ if (!PendingNodes.empty() && !PendingRequest) {
+ const auto& nodeId = PendingNodes.front();
+ auto kqpProxyId = NKqp::MakeKqpProxyID(nodeId);
+ auto req = std::make_unique<NKikimr::NKqp::TEvKqp::TEvListSessionsRequest>();
+ if (!ContinuationToken.empty()) {
+ req->Record.SetSessionIdStart(ContinuationToken);
+ req->Record.SetSessionIdStartInclusive(true);
+ } else {
+ req->Record.SetSessionIdStart(SessionIdFrom);
+ req->Record.SetSessionIdStartInclusive(SessionIdFromInclusive);
+ }
+
+ req->Record.SetSessionIdEnd(SessionIdTo);
+ req->Record.SetSessionIdEndInclusive(SessionIdToInclusive);
+ req->Record.MutableColumns()->Add(ColumnsToRead.begin(), ColumnsToRead.end());
+ if (FreeSpace == 0) {
+ FreeSpace = 1_KB;
+ }
+
+ req->Record.SetFreeSpace(FreeSpace);
+
+ LOG_DEBUG_S(TlsActivationContext->AsActorContext(), NKikimrServices::SYSTEM_VIEWS,
+ "Send request to node, node_id=" << nodeId << ", request: " << req->Record.ShortDebugString());
+
+ Send(kqpProxyId, req.release(), 0, nodeId);
+ PendingRequest = true;
+ }
+ }
+
+ void Handle(NKqp::TEvKqp::TEvListProxyNodesResponse::TPtr& ev) {
+ auto& proxies = ev->Get()->ProxyNodes;
+ std::sort(proxies.begin(), proxies.end());
+ PendingNodes = std::deque<ui32>(proxies.begin(), proxies.end());
+ PendingNodesInitialized = true;
+ StartScan();
+ }
+
+ void Handle(NKqp::TEvKqpCompute::TEvScanDataAck::TPtr& ev) {
+ FreeSpace = ev->Get()->FreeSpace;
+ StartScan();
+ }
+
+ void Handle(NKqp::TEvKqp::TEvListSessionsResponse::TPtr& ev) {
+ auto& record = ev->Get()->Record;
+ LastResponse = std::move(record);
+ ProcessRows();
+ }
+
+ void Undelivered(TEvents::TEvUndelivered::TPtr& ev) {
+ if (ev->Get()->SourceType == NKqp::TKqpEvents::EvListSessionsRequest) {
+ ui32 nodeId = ev->Cookie;
+ LOG_INFO_S(TlsActivationContext->AsActorContext(), NKikimrServices::SYSTEM_VIEWS,
+ "Received undelivered response for node_id: " << nodeId);
+ StartScan();
+ }
+ }
+
+ void Connected(TEvInterconnect::TEvNodeConnected::TPtr&) {
+ }
+
+ void Disconnected(TEvInterconnect::TEvNodeDisconnected::TPtr& ev) {
+ ui32 nodeId = ev->Get()->NodeId;
+ Y_UNUSED(nodeId);
+ ProcessRows();
+ }
+
+ void ProcessRows() {
+ auto batch = MakeHolder<NKqp::TEvKqpCompute::TEvScanData>(ScanId);
+ auto nodeId = LastResponse.GetNodeId();
+ for(int idx = 0; idx < LastResponse.GetSessions().size(); ++idx) {
+ TVector<TCell> cells;
+ for (auto extractor : ColumnsExtractors) {
+ cells.push_back(extractor(LastResponse.GetSessions(idx), nodeId));
+ }
+
+ TArrayRef<const TCell> ref(cells);
+ batch->Rows.emplace_back(TOwnedCellVec::Make(ref));
+ cells.clear();
+ }
+
+ if (LastResponse.GetFinished()) {
+ PendingNodes.pop_front();
+ ContinuationToken = TString();
+ if (PendingNodes.empty()) {
+ batch->Finished = true;
+ }
+
+ } else {
+ ContinuationToken = LastResponse.GetContinuationToken();
+ }
+
+ PendingRequest = false;
+ SendBatch(std::move(batch));
+ StartScan();
+ }
+
+ void PassAway() override {
+ TBase::PassAway();
+ }
+
+private:
+ TString SessionIdFrom;
+ bool SessionIdFromInclusive = false;
+ TString SessionIdTo;
+ bool SessionIdToInclusive = false;
+
+ TString ContinuationToken;
+
+ bool PendingRequest = false;
+ bool IsEmptyRange = false;
+ std::vector<ui32> MissingSchemaColumns;
+ std::deque<ui32> PendingNodes;
+ i64 FreeSpace = 0;
+ bool PendingNodesInitialized = false;
+ std::vector<TExtractor> ColumnsExtractors;
+ std::vector<ui32> ColumnsToRead;
+ NKikimrKqp::TEvListSessionsResponse LastResponse;
+};
+
+THolder<NActors::IActor> CreateSessionsScan(const NActors::TActorId& ownerId, ui32 scanId, const TTableId& tableId,
+ const TTableRange& tableRange, const TArrayRef<NMiniKQL::TKqpComputeContextBase::TColumn>& columns)
+{
+ return MakeHolder<TSessionsScan>(ownerId, scanId, tableId, tableRange, columns);
+}
+
+} // NKikimr::NSysView
diff --git a/ydb/core/sys_view/sessions/sessions.h b/ydb/core/sys_view/sessions/sessions.h
new file mode 100644
index 00000000000..6384626d195
--- /dev/null
+++ b/ydb/core/sys_view/sessions/sessions.h
@@ -0,0 +1,15 @@
+#pragma once
+
+#include <ydb/core/kqp/runtime/kqp_compute.h>
+
+#include <ydb/library/actors/core/actor.h>
+#include <ydb/library/actors/core/actorid.h>
+
+namespace NKikimr {
+namespace NSysView {
+
+THolder<NActors::IActor> CreateSessionsScan(const NActors::TActorId& ownerId, ui32 scanId, const TTableId& tableId,
+ const TTableRange& tableRange, const TArrayRef<NMiniKQL::TKqpComputeContextBase::TColumn>& columns);
+
+} // NSysView
+} // NKikimr
diff --git a/ydb/core/sys_view/sessions/ya.make b/ydb/core/sys_view/sessions/ya.make
new file mode 100644
index 00000000000..0c1ff7dd4f4
--- /dev/null
+++ b/ydb/core/sys_view/sessions/ya.make
@@ -0,0 +1,21 @@
+LIBRARY()
+
+OWNER(
+ g:kikimr
+)
+
+SRCS(
+ sessions.h
+ sessions.cpp
+)
+
+PEERDIR(
+ ydb/library/actors/core
+ ydb/core/base
+ ydb/core/kqp/runtime
+ ydb/core/sys_view/common
+)
+
+YQL_LAST_ABI_VERSION()
+
+END()
diff --git a/ydb/core/sys_view/ut_kqp.cpp b/ydb/core/sys_view/ut_kqp.cpp
index bd930aa7040..91ffe34776a 100644
--- a/ydb/core/sys_view/ut_kqp.cpp
+++ b/ydb/core/sys_view/ut_kqp.cpp
@@ -1526,7 +1526,7 @@ Y_UNIT_TEST_SUITE(SystemView) {
UNIT_ASSERT_VALUES_EQUAL(entry.Type, ESchemeEntryType::Directory);
auto children = result.GetChildren();
- UNIT_ASSERT_VALUES_EQUAL(children.size(), 18);
+ UNIT_ASSERT_VALUES_EQUAL(children.size(), 19);
THashSet<TString> names;
for (const auto& child : children) {
@@ -1544,7 +1544,7 @@ Y_UNIT_TEST_SUITE(SystemView) {
UNIT_ASSERT_VALUES_EQUAL(entry.Type, ESchemeEntryType::Directory);
auto children = result.GetChildren();
- UNIT_ASSERT_VALUES_EQUAL(children.size(), 12);
+ UNIT_ASSERT_VALUES_EQUAL(children.size(), 13);
THashSet<TString> names;
for (const auto& child : children) {
diff --git a/ydb/core/sys_view/ya.make b/ydb/core/sys_view/ya.make
index 057a42fc198..0fa347cf380 100644
--- a/ydb/core/sys_view/ya.make
+++ b/ydb/core/sys_view/ya.make
@@ -10,6 +10,7 @@ PEERDIR(
ydb/core/kqp/runtime
ydb/core/sys_view/common
ydb/core/sys_view/nodes
+ ydb/core/sys_view/sessions
ydb/core/sys_view/partition_stats
ydb/core/sys_view/query_stats
ydb/core/sys_view/service