diff options
| author | azevaykin <[email protected]> | 2024-07-30 15:44:31 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-07-30 15:44:31 +0300 |
| commit | 5ead379fac334f64891a48dbb82be3f64a8b02ca (patch) | |
| tree | dc5f0a6cdd4fb7cc5029dcccb4cb4d1b3aec1992 | |
| parent | e661cdaf98bb1e1242020610c34d5edce5ad85d5 (diff) | |
Refactoring statistics: scan -> traversal (#7210)
19 files changed, 403 insertions, 402 deletions
diff --git a/ydb/core/protos/counters_statistics_aggregator.proto b/ydb/core/protos/counters_statistics_aggregator.proto index e6bf4adc9a3..64f4cc4c63d 100644 --- a/ydb/core/protos/counters_statistics_aggregator.proto +++ b/ydb/core/protos/counters_statistics_aggregator.proto @@ -11,12 +11,12 @@ enum ETxTypes { TXTYPE_INIT = 1 [(TxTypeOpts) = {Name: "TxInit"}]; TXTYPE_CONFIGURE = 2 [(TxTypeOpts) = {Name: "TxConfigure"}]; TXTYPE_SCHEMESHARD_STATS = 3 [(TxTypeOpts) = {Name: "TxSchemeShardStats"}]; - TXTYPE_SCAN_TABLE = 4 [(TxTypeOpts) = {Name: "TxScanTable"}]; + TXTYPE_ANALYZE_TABLE = 4 [(TxTypeOpts) = {Name: "TxAnalyzeTable"}]; TXTYPE_NAVIGATE = 5 [(TxTypeOpts) = {Name: "TxNavigate"}]; TXTYPE_RESOLVE = 6 [(TxTypeOpts) = {Name: "TxResolve"}]; TXTYPE_SCAN_RESPONSE = 7 [(TxTypeOpts) = {Name: "TxScanResponse"}]; TXTYPE_SAVE_QUERY_RESPONSE = 8 [(TxTypeOpts) = {Name: "TxSaveQueryResponse"}]; - TXTYPE_SCHEDULE_SCAN = 9 [(TxTypeOpts) = {Name: "TxScheduleScan"}]; + TXTYPE_SCHEDULE_TRAVERSAL = 9 [(TxTypeOpts) = {Name: "TxScheduleTraversal"}]; TXTYPE_DELETE_QUERY_RESPONSE = 10 [(TxTypeOpts) = {Name: "TxDeleteQueryResponse"}]; TXTYPE_AGGR_STAT_RESPONSE = 11 [(TxTypeOpts) = {Name: "TxAggregateStatisticsResponse"}]; TXTYPE_RESPONSE_TABLET_DISTRIBUTION = 12 [(TxTypeOpts) = {Name: "TxResponseTabletDistribution"}]; diff --git a/ydb/core/statistics/aggregator/aggregator_impl.cpp b/ydb/core/statistics/aggregator/aggregator_impl.cpp index 110ef62d59d..3bdea2e2802 100644 --- a/ydb/core/statistics/aggregator/aggregator_impl.cpp +++ b/ydb/core/statistics/aggregator/aggregator_impl.cpp @@ -383,8 +383,8 @@ size_t TStatisticsAggregator::PropagatePart(const std::vector<TNodeId>& nodeIds, auto ssId = ssIds[index]; auto* entry = record->AddEntries(); entry->SetSchemeShardId(ssId); - auto itStats = BaseStats.find(ssId); - if (itStats != BaseStats.end()) { + auto itStats = BaseStatistics.find(ssId); + if (itStats != BaseStatistics.end()) { entry->SetStats(itStats->second); size += itStats->second.size(); } else { @@ -398,21 +398,21 @@ size_t TStatisticsAggregator::PropagatePart(const std::vector<TNodeId>& nodeIds, } void TStatisticsAggregator::Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - if (!ScanTableId.PathId) { + if (!TraversalTableId.PathId) { return; } auto tabletId = ev->Get()->TabletId; - if (IsColumnTable) { + if (TraversalIsColumnTable) { if (tabletId == HiveId) { Schedule(HiveRetryInterval, new TEvPrivate::TEvRequestDistribution); } else { SA_LOG_CRIT("[" << TabletID() << "] TEvDeliveryProblem with unexpected tablet " << tabletId); } } else { - if (ShardRanges.empty()) { + if (DatashardRanges.empty()) { return; } - auto& range = ShardRanges.front(); + auto& range = DatashardRanges.front(); if (tabletId != range.DataShardId) { return; } @@ -439,11 +439,11 @@ void TStatisticsAggregator::Handle(TEvStatistics::TEvAnalyzeStatus::TPtr& ev) { auto response = std::make_unique<TEvStatistics::TEvAnalyzeStatusResponse>(); auto& outRecord = response->Record; - if (ScanTableId.PathId == pathId) { + if (TraversalTableId.PathId == pathId) { outRecord.SetStatus(NKikimrStat::TEvAnalyzeStatusResponse::STATUS_IN_PROGRESS); } else { - auto it = ScanOperationsByPathId.find(pathId); - if (it != ScanOperationsByPathId.end()) { + auto it = ForceTraversalsByPathId.find(pathId); + if (it != ForceTraversalsByPathId.end()) { outRecord.SetStatus(NKikimrStat::TEvAnalyzeStatusResponse::STATUS_ENQUEUED); } else { outRecord.SetStatus(NKikimrStat::TEvAnalyzeStatusResponse::STATUS_NO_OPERATION); @@ -486,7 +486,7 @@ void TStatisticsAggregator::InitializeStatisticsTable() { void TStatisticsAggregator::Navigate() { using TNavigate = NSchemeCache::TSchemeCacheNavigate; TNavigate::TEntry entry; - entry.TableId = ScanTableId; + entry.TableId = TraversalTableId; entry.RequestType = TNavigate::TEntry::ERequestType::ByTableId; entry.Operation = TNavigate::OpTable; @@ -500,9 +500,9 @@ void TStatisticsAggregator::Resolve() { ++ResolveRound; TVector<TCell> plusInf; - TTableRange range(StartKey.GetCells(), true, plusInf, true, false); + TTableRange range(TraversalStartKey.GetCells(), true, plusInf, true, false); auto keyDesc = MakeHolder<TKeyDesc>( - ScanTableId, range, TKeyDesc::ERowOperation::Read, KeyColumnTypes, Columns); + TraversalTableId, range, TKeyDesc::ERowOperation::Read, KeyColumnTypes, Columns); auto request = std::make_unique<NSchemeCache::TSchemeCacheRequest>(); request->ResultSet.emplace_back(std::move(keyDesc)); @@ -510,19 +510,19 @@ void TStatisticsAggregator::Resolve() { Send(MakeSchemeCacheID(), new TEvTxProxySchemeCache::TEvResolveKeySet(request.release())); } -void TStatisticsAggregator::NextRange() { - if (ShardRanges.empty()) { +void TStatisticsAggregator::ScanNextDatashardRange() { + if (DatashardRanges.empty()) { SaveStatisticsToTable(); return; } - auto& range = ShardRanges.front(); + auto& range = DatashardRanges.front(); auto request = std::make_unique<NStat::TEvStatistics::TEvStatisticsRequest>(); auto& record = request->Record; auto* path = record.MutableTable()->MutablePathId(); - path->SetOwnerId(ScanTableId.PathId.OwnerId); - path->SetLocalId(ScanTableId.PathId.LocalPathId); - record.SetStartKey(StartKey.GetBuffer()); + path->SetOwnerId(TraversalTableId.PathId.OwnerId); + path->SetLocalId(TraversalTableId.PathId.LocalPathId); + record.SetStartKey(TraversalStartKey.GetBuffer()); Send(MakePipePerNodeCacheID(false), new TEvPipeCache::TEvForward(request.release(), range.DataShardId, true), @@ -556,7 +556,7 @@ void TStatisticsAggregator::SaveStatisticsToTable() { data.push_back(strSketch); } - Register(CreateSaveStatisticsQuery(ScanTableId.PathId, EStatType::COUNT_MIN_SKETCH, + Register(CreateSaveStatisticsQuery(TraversalTableId.PathId, EStatType::COUNT_MIN_SKETCH, std::move(columnTags), std::move(data))); } @@ -568,76 +568,76 @@ void TStatisticsAggregator::DeleteStatisticsFromTable() { PendingDeleteStatistics = false; - Register(CreateDeleteStatisticsQuery(ScanTableId.PathId)); + Register(CreateDeleteStatisticsQuery(TraversalTableId.PathId)); } -void TStatisticsAggregator::ScheduleNextScan(NIceDb::TNiceDb& db) { - if (!ScanOperations.Empty()) { - auto* operation = ScanOperations.Front(); +void TStatisticsAggregator::ScheduleNextTraversal(NIceDb::TNiceDb& db) { + if (!ForceTraversals.Empty()) { + auto* operation = ForceTraversals.Front(); ReplyToActorIds.swap(operation->ReplyToActorIds); - bool doStartScan = true; + bool doStartAnalyze = true; bool isColumnTable = false; auto pathId = operation->PathId; - auto itPath = ScanTables.find(pathId); - if (itPath != ScanTables.end()) { + auto itPath = ScheduleTraversals.find(pathId); + if (itPath != ScheduleTraversals.end()) { isColumnTable = itPath->second.IsColumnTable; } else { - doStartScan = false; + doStartAnalyze = false; } - if (doStartScan) { - StartScan(db, pathId, isColumnTable); + if (doStartAnalyze) { + StartTraversal(db, pathId, isColumnTable); } - db.Table<Schema::ScanOperations>().Key(operation->OperationId).Delete(); - ScanOperations.PopFront(); - ScanOperationsByPathId.erase(pathId); - return; - } - if (ScanTablesByTime.Empty()) { - return; - } - auto* topTable = ScanTablesByTime.Top(); - if (TInstant::Now() < topTable->LastUpdateTime + ScanIntervalTime) { - return; - } - bool isColumnTable = false; - auto itPath = ScanTables.find(topTable->PathId); - if (itPath != ScanTables.end()) { - isColumnTable = itPath->second.IsColumnTable; - } else { - return; + db.Table<Schema::ForceTraversals>().Key(operation->OperationId).Delete(); + ForceTraversals.PopFront(); + ForceTraversalsByPathId.erase(pathId); + } else { // ForceTraversals is empty, then go to ScheduleTraversals + if (ScheduleTraversalsByTime.Empty()) { + return; + } + auto* oldestTable = ScheduleTraversalsByTime.Top(); + if (TInstant::Now() < oldestTable->LastUpdateTime + ScheduleTraversalPeriod) { + return; + } + bool isColumnTable = false; + auto itPath = ScheduleTraversals.find(oldestTable->PathId); + if (itPath != ScheduleTraversals.end()) { + isColumnTable = itPath->second.IsColumnTable; + } else { + return; + } + StartTraversal(db, oldestTable->PathId, isColumnTable); } - StartScan(db, topTable->PathId, isColumnTable); } -void TStatisticsAggregator::StartScan(NIceDb::TNiceDb& db, TPathId pathId, bool isColumnTable) { - ScanTableId.PathId = pathId; - ScanStartTime = TInstant::Now(); - IsColumnTable = isColumnTable; - PersistCurrentScan(db); +void TStatisticsAggregator::StartTraversal(NIceDb::TNiceDb& db, TPathId pathId, bool isColumnTable) { + TraversalTableId.PathId = pathId; + TraversalStartTime = TInstant::Now(); + TraversalIsColumnTable = isColumnTable; + PersistTraversal(db); - StartKey = TSerializedCellVec(); + TraversalStartKey = TSerializedCellVec(); PersistStartKey(db); Navigate(); } -void TStatisticsAggregator::FinishScan(NIceDb::TNiceDb& db) { - auto pathId = ScanTableId.PathId; +void TStatisticsAggregator::FinishTraversal(NIceDb::TNiceDb& db) { + auto pathId = TraversalTableId.PathId; - auto pathIt = ScanTables.find(pathId); - if (pathIt != ScanTables.end()) { - auto& scanTable = pathIt->second; - scanTable.LastUpdateTime = ScanStartTime; - db.Table<Schema::ScanTables>().Key(pathId.OwnerId, pathId.LocalPathId).Update( - NIceDb::TUpdate<Schema::ScanTables::LastUpdateTime>(ScanStartTime.MicroSeconds())); + auto pathIt = ScheduleTraversals.find(pathId); + if (pathIt != ScheduleTraversals.end()) { + auto& traversalTable = pathIt->second; + traversalTable.LastUpdateTime = TraversalStartTime; + db.Table<Schema::ScheduleTraversals>().Key(pathId.OwnerId, pathId.LocalPathId).Update( + NIceDb::TUpdate<Schema::ScheduleTraversals::LastUpdateTime>(TraversalStartTime.MicroSeconds())); - if (ScanTablesByTime.Has(&scanTable)) { - ScanTablesByTime.Update(&scanTable); + if (ScheduleTraversalsByTime.Has(&traversalTable)) { + ScheduleTraversalsByTime.Update(&traversalTable); } } - ResetScanState(db); + ResetTraversalState(db); } void TStatisticsAggregator::PersistSysParam(NIceDb::TNiceDb& db, ui64 id, const TString& value) { @@ -645,41 +645,41 @@ void TStatisticsAggregator::PersistSysParam(NIceDb::TNiceDb& db, ui64 id, const NIceDb::TUpdate<Schema::SysParams::Value>(value)); } -void TStatisticsAggregator::PersistCurrentScan(NIceDb::TNiceDb& db) { - PersistSysParam(db, Schema::SysParam_ScanTableOwnerId, ToString(ScanTableId.PathId.OwnerId)); - PersistSysParam(db, Schema::SysParam_ScanTableLocalPathId, ToString(ScanTableId.PathId.LocalPathId)); - PersistSysParam(db, Schema::SysParam_ScanStartTime, ToString(ScanStartTime.MicroSeconds())); - PersistSysParam(db, Schema::SysParam_IsColumnTable, ToString(IsColumnTable)); +void TStatisticsAggregator::PersistTraversal(NIceDb::TNiceDb& db) { + PersistSysParam(db, Schema::SysParam_TraversalTableOwnerId, ToString(TraversalTableId.PathId.OwnerId)); + PersistSysParam(db, Schema::SysParam_TraversalTableLocalPathId, ToString(TraversalTableId.PathId.LocalPathId)); + PersistSysParam(db, Schema::SysParam_TraversalStartTime, ToString(TraversalStartTime.MicroSeconds())); + PersistSysParam(db, Schema::SysParam_TraversalIsColumnTable, ToString(TraversalIsColumnTable)); } void TStatisticsAggregator::PersistStartKey(NIceDb::TNiceDb& db) { - PersistSysParam(db, Schema::SysParam_StartKey, StartKey.GetBuffer()); + PersistSysParam(db, Schema::SysParam_TraversalStartKey, TraversalStartKey.GetBuffer()); } -void TStatisticsAggregator::PersistLastScanOperationId(NIceDb::TNiceDb& db) { - PersistSysParam(db, Schema::SysParam_LastScanOperationId, ToString(LastScanOperationId)); +void TStatisticsAggregator::PersistLastForceTraversalOperationId(NIceDb::TNiceDb& db) { + PersistSysParam(db, Schema::SysParam_LastForceTraversalOperationId, ToString(LastForceTraversalOperationId)); } void TStatisticsAggregator::PersistGlobalTraversalRound(NIceDb::TNiceDb& db) { PersistSysParam(db, Schema::SysParam_GlobalTraversalRound, ToString(GlobalTraversalRound)); } -void TStatisticsAggregator::ResetScanState(NIceDb::TNiceDb& db) { - ScanTableId.PathId = TPathId(); - ScanStartTime = TInstant::MicroSeconds(0); - PersistCurrentScan(db); +void TStatisticsAggregator::ResetTraversalState(NIceDb::TNiceDb& db) { + TraversalTableId.PathId = TPathId(); + TraversalStartTime = TInstant::MicroSeconds(0); + PersistTraversal(db); - StartKey = TSerializedCellVec(); + TraversalStartKey = TSerializedCellVec(); PersistStartKey(db); ReplyToActorIds.clear(); for (auto& [tag, _] : CountMinSketches) { - db.Table<Schema::Statistics>().Key(tag).Delete(); + db.Table<Schema::ColumnStatistics>().Key(tag).Delete(); } CountMinSketches.clear(); - ShardRanges.clear(); + DatashardRanges.clear(); KeyColumnTypes.clear(); Columns.clear(); @@ -728,7 +728,7 @@ bool TStatisticsAggregator::OnRenderAppHtmlPage(NMon::TEvRemoteHttpInfo::TPtr ev PRE() { str << "---- StatisticsAggregator ----" << Endl << Endl; str << "Database: " << Database << Endl; - str << "BaseStats: " << BaseStats.size() << Endl; + str << "BaseStatistics: " << BaseStatistics.size() << Endl; str << "SchemeShards: " << SchemeShards.size() << Endl; { std::function<TSSId(const std::pair<const TSSId, size_t>&)> extr = @@ -773,24 +773,24 @@ bool TStatisticsAggregator::OnRenderAppHtmlPage(NMon::TEvRemoteHttpInfo::TPtr ev str << "PendingRequests: " << PendingRequests.size() << Endl; str << "ProcessUrgentInFlight: " << ProcessUrgentInFlight << Endl << Endl; - str << "ScanTableId: " << ScanTableId << Endl; + str << "TraversalTableId: " << TraversalTableId << Endl; str << "Columns: " << Columns.size() << Endl; - str << "ShardRanges: " << ShardRanges.size() << Endl; + str << "DatashardRanges: " << DatashardRanges.size() << Endl; str << "CountMinSketches: " << CountMinSketches.size() << Endl << Endl; - str << "ScanTablesByTime: " << ScanTablesByTime.Size() << Endl; - if (!ScanTablesByTime.Empty()) { - auto* scanTable = ScanTablesByTime.Top(); - str << " top: " << scanTable->PathId - << ", last update time: " << scanTable->LastUpdateTime << Endl; + str << "ScheduleTraversalsByTime: " << ScheduleTraversalsByTime.Size() << Endl; + if (!ScheduleTraversalsByTime.Empty()) { + auto* oldestTable = ScheduleTraversalsByTime.Top(); + str << " oldest table: " << oldestTable->PathId + << ", ordest table update time: " << oldestTable->LastUpdateTime << Endl; } - str << "ScanTablesBySchemeShard: " << ScanTablesBySchemeShard.size() << Endl; - if (!ScanTablesBySchemeShard.empty()) { - str << " " << ScanTablesBySchemeShard.begin()->first << Endl; + str << "ScheduleTraversalsBySchemeShard: " << ScheduleTraversalsBySchemeShard.size() << Endl; + if (!ScheduleTraversalsBySchemeShard.empty()) { + str << " " << ScheduleTraversalsBySchemeShard.begin()->first << Endl; std::function<TPathId(const TPathId&)> extr = [](const auto& x) { return x; }; - PrintContainerStart(ScanTablesBySchemeShard.begin()->second, 2, str, extr); + PrintContainerStart(ScheduleTraversalsBySchemeShard.begin()->second, 2, str, extr); } - str << "ScanStartTime: " << ScanStartTime << Endl; + str << "TraversalStartTime: " << TraversalStartTime << Endl; } } diff --git a/ydb/core/statistics/aggregator/aggregator_impl.h b/ydb/core/statistics/aggregator/aggregator_impl.h index 333109e297f..a6dd1c4ead8 100644 --- a/ydb/core/statistics/aggregator/aggregator_impl.h +++ b/ydb/core/statistics/aggregator/aggregator_impl.h @@ -45,12 +45,12 @@ private: struct TTxInit; struct TTxConfigure; struct TTxSchemeShardStats; - struct TTxScanTable; + struct TTxAnalyzeTable; struct TTxNavigate; struct TTxResolve; - struct TTxStatisticsScanResponse; + struct TTxDatashardScanResponse; struct TTxSaveQueryResponse; - struct TTxScheduleScan; + struct TTxScheduleTrasersal; struct TTxDeleteQueryResponse; struct TTxAggregateStatisticsResponse; struct TTxResponseTabletDistribution; @@ -62,7 +62,7 @@ private: EvFastPropagateCheck, EvProcessUrgent, EvPropagateTimeout, - EvScheduleScan, + EvScheduleTraversal, EvRequestDistribution, EvResolve, EvAckTimeout, @@ -74,7 +74,7 @@ private: struct TEvFastPropagateCheck : public TEventLocal<TEvFastPropagateCheck, EvFastPropagateCheck> {}; struct TEvProcessUrgent : public TEventLocal<TEvProcessUrgent, EvProcessUrgent> {}; struct TEvPropagateTimeout : public TEventLocal<TEvPropagateTimeout, EvPropagateTimeout> {}; - struct TEvScheduleScan : public TEventLocal<TEvScheduleScan, EvScheduleScan> {}; + struct TEvScheduleTraversal : public TEventLocal<TEvScheduleTraversal, EvScheduleTraversal> {}; struct TEvRequestDistribution : public TEventLocal<TEvRequestDistribution, EvRequestDistribution> {}; struct TEvResolve : public TEventLocal<TEvResolve, EvResolve> {}; @@ -128,7 +128,7 @@ private: void Handle(TEvStatistics::TEvStatTableCreationResponse::TPtr& ev); void Handle(TEvStatistics::TEvSaveStatisticsQueryResponse::TPtr& ev); void Handle(TEvStatistics::TEvDeleteStatisticsQueryResponse::TPtr& ev); - void Handle(TEvPrivate::TEvScheduleScan::TPtr& ev); + void Handle(TEvPrivate::TEvScheduleTraversal::TPtr& ev); void Handle(TEvStatistics::TEvAnalyzeStatus::TPtr& ev); void Handle(TEvHive::TEvResponseTabletDistribution::TPtr& ev); void Handle(TEvStatistics::TEvAggregateStatisticsResponse::TPtr& ev); @@ -140,20 +140,20 @@ private: void InitializeStatisticsTable(); void Navigate(); void Resolve(); - void NextRange(); + void ScanNextDatashardRange(); void SaveStatisticsToTable(); void DeleteStatisticsFromTable(); void PersistSysParam(NIceDb::TNiceDb& db, ui64 id, const TString& value); - void PersistCurrentScan(NIceDb::TNiceDb& db); + void PersistTraversal(NIceDb::TNiceDb& db); void PersistStartKey(NIceDb::TNiceDb& db); - void PersistLastScanOperationId(NIceDb::TNiceDb& db); + void PersistLastForceTraversalOperationId(NIceDb::TNiceDb& db); void PersistGlobalTraversalRound(NIceDb::TNiceDb& db); - void ResetScanState(NIceDb::TNiceDb& db); - void ScheduleNextScan(NIceDb::TNiceDb& db); - void StartScan(NIceDb::TNiceDb& db, TPathId pathId, bool isColumnTable); - void FinishScan(NIceDb::TNiceDb& db); + void ResetTraversalState(NIceDb::TNiceDb& db); + void ScheduleNextTraversal(NIceDb::TNiceDb& db); + void StartTraversal(NIceDb::TNiceDb& db, TPathId pathId, bool isColumnTable); + void FinishTraversal(NIceDb::TNiceDb& db); STFUNC(StateInit) { StateInitImpl(ev, SelfId()); @@ -184,7 +184,7 @@ private: hFunc(TEvStatistics::TEvStatTableCreationResponse, Handle); hFunc(TEvStatistics::TEvSaveStatisticsQueryResponse, Handle); hFunc(TEvStatistics::TEvDeleteStatisticsQueryResponse, Handle); - hFunc(TEvPrivate::TEvScheduleScan, Handle); + hFunc(TEvPrivate::TEvScheduleTraversal, Handle); hFunc(TEvStatistics::TEvAnalyzeStatus, Handle); hFunc(TEvHive::TEvResponseTabletDistribution, Handle); hFunc(TEvStatistics::TEvAggregateStatisticsResponse, Handle); @@ -216,7 +216,7 @@ private: TDuration PropagateTimeout; static constexpr TDuration FastCheckInterval = TDuration::MilliSeconds(50); - std::unordered_map<TSSId, TString> BaseStats; // schemeshard id -> serialized stats for all paths + std::unordered_map<TSSId, TString> BaseStatistics; // schemeshard id -> serialized stats for all paths std::unordered_map<TSSId, size_t> SchemeShards; // all connected schemeshards std::unordered_map<TActorId, TSSId> SchemeShardPipes; // schemeshard pipe servers @@ -239,10 +239,6 @@ private: std::queue<TEvStatistics::TEvRequestStats::TPtr> PendingRequests; bool ProcessUrgentInFlight = false; - // - - TTableId ScanTableId; // stored in local db - bool IsColumnTable = false; // stored in local db std::unordered_set<TActorId> ReplyToActorIds; bool IsStatisticsTableCreated = false; @@ -257,16 +253,14 @@ private: TSerializedCellVec EndKey; ui64 DataShardId = 0; }; - std::deque<TRange> ShardRanges; - - TSerializedCellVec StartKey; // stored in local db - - std::unordered_map<ui32, std::unique_ptr<TCountMinSketch>> CountMinSketches; // stored in local db + std::deque<TRange> DatashardRanges; - static constexpr TDuration ScanIntervalTime = TDuration::Hours(24); - static constexpr TDuration ScheduleScanIntervalTime = TDuration::Seconds(1); + // period for both force and schedule traversals + static constexpr TDuration TraversalPeriod = TDuration::Seconds(1); + // if table traverse time is older, than traserse it on schedule + static constexpr TDuration ScheduleTraversalPeriod = TDuration::Hours(24); - struct TScanTable { + struct TScheduleTraversal { TPathId PathId; ui64 SchemeShardId = 0; TInstant LastUpdateTime; @@ -275,36 +269,17 @@ private: size_t HeapIndexByTime = -1; struct THeapIndexByTime { - size_t& operator()(TScanTable& value) const { + size_t& operator()(TScheduleTraversal& value) const { return value.HeapIndexByTime; } }; struct TLessByTime { - bool operator()(const TScanTable& l, const TScanTable& r) const { + bool operator()(const TScheduleTraversal& l, const TScheduleTraversal& r) const { return l.LastUpdateTime < r.LastUpdateTime; } }; }; - std::unordered_map<TPathId, TScanTable> ScanTables; // stored in local db - - std::unordered_map<ui64, std::unordered_set<TPathId>> ScanTablesBySchemeShard; - - typedef TIntrusiveHeap<TScanTable, TScanTable::THeapIndexByTime, TScanTable::TLessByTime> - TScanTableQueueByTime; - TScanTableQueueByTime ScanTablesByTime; - - struct TScanOperation : public TIntrusiveListItem<TScanOperation> { - ui64 OperationId = 0; - TPathId PathId; - std::unordered_set<TActorId> ReplyToActorIds; - }; - TIntrusiveList<TScanOperation> ScanOperations; // stored in local db - std::unordered_map<TPathId, TScanOperation> ScanOperationsByPathId; - - ui64 LastScanOperationId = 0; // stored in local db - - TInstant ScanStartTime; size_t ResolveRound = 0; static constexpr size_t MaxResolveRoundCount = 5; @@ -319,10 +294,36 @@ private: size_t TraversalRound = 0; static constexpr size_t MaxTraversalRoundCount = 5; - size_t GlobalTraversalRound = 1; // stored in local db + size_t KeepAliveSeqNo = 0; static constexpr TDuration KeepAliveTimeout = TDuration::Seconds(3); + +private: // stored in local db + + TTableId TraversalTableId; + bool TraversalIsColumnTable = false; + TSerializedCellVec TraversalStartKey; + TInstant TraversalStartTime; + ui64 LastForceTraversalOperationId = 0; + + size_t GlobalTraversalRound = 1; + + std::unordered_map<ui32, std::unique_ptr<TCountMinSketch>> CountMinSketches; + + std::unordered_map<TPathId, TScheduleTraversal> ScheduleTraversals; + std::unordered_map<ui64, std::unordered_set<TPathId>> ScheduleTraversalsBySchemeShard; + typedef TIntrusiveHeap<TScheduleTraversal, TScheduleTraversal::THeapIndexByTime, TScheduleTraversal::TLessByTime> + TTraversalsByTime; + TTraversalsByTime ScheduleTraversalsByTime; + + struct TForceTraversal : public TIntrusiveListItem<TForceTraversal> { + ui64 OperationId = 0; + TPathId PathId; + std::unordered_set<TActorId> ReplyToActorIds; + }; + TIntrusiveList<TForceTraversal> ForceTraversals; + std::unordered_map<TPathId, TForceTraversal> ForceTraversalsByPathId; }; } // NKikimr::NStat diff --git a/ydb/core/statistics/aggregator/schema.h b/ydb/core/statistics/aggregator/schema.h index e8751204568..617725ed974 100644 --- a/ydb/core/statistics/aggregator/schema.h +++ b/ydb/core/statistics/aggregator/schema.h @@ -13,7 +13,7 @@ struct TAggregatorSchema : NIceDb::Schema { using TColumns = TableColumns<Id, Value>; }; - struct BaseStats : Table<2> { + struct BaseStatistics : Table<2> { struct SchemeShardId : Column<1, NScheme::NTypeIds::Uint64> {}; struct Stats : Column<2, NScheme::NTypeIds::String> {}; @@ -21,7 +21,7 @@ struct TAggregatorSchema : NIceDb::Schema { using TColumns = TableColumns<SchemeShardId, Stats>; }; - struct Statistics : Table<3> { + struct ColumnStatistics : Table<3> { struct ColumnTag : Column<1, NScheme::NTypeIds::Uint32> {}; struct CountMinSketch : Column<2, NScheme::NTypeIds::String> {}; @@ -29,7 +29,7 @@ struct TAggregatorSchema : NIceDb::Schema { using TColumns = TableColumns<ColumnTag, CountMinSketch>; }; - struct ScanTables : Table<4> { + struct ScheduleTraversals : Table<4> { struct OwnerId : Column<1, NScheme::NTypeIds::Uint64> {}; struct LocalPathId : Column<2, NScheme::NTypeIds::Uint64> {}; struct LastUpdateTime : Column<3, NScheme::NTypeIds::Timestamp> {}; @@ -46,7 +46,7 @@ struct TAggregatorSchema : NIceDb::Schema { >; }; - struct ScanOperations : Table<5> { + struct ForceTraversals : Table<5> { struct OperationId : Column<1, NScheme::NTypeIds::Uint64> {}; struct OwnerId : Column<2, NScheme::NTypeIds::Uint64> {}; struct LocalPathId : Column<3, NScheme::NTypeIds::Uint64> {}; @@ -61,10 +61,10 @@ struct TAggregatorSchema : NIceDb::Schema { using TTables = SchemaTables< SysParams, - BaseStats, - Statistics, - ScanTables, - ScanOperations + BaseStatistics, + ColumnStatistics, + ScheduleTraversals, + ForceTraversals >; using TSettings = SchemaSettings< @@ -73,12 +73,12 @@ struct TAggregatorSchema : NIceDb::Schema { >; static constexpr ui64 SysParam_Database = 1; - static constexpr ui64 SysParam_StartKey = 2; - static constexpr ui64 SysParam_ScanTableOwnerId = 3; - static constexpr ui64 SysParam_ScanTableLocalPathId = 4; - static constexpr ui64 SysParam_ScanStartTime = 5; - static constexpr ui64 SysParam_LastScanOperationId = 6; - static constexpr ui64 SysParam_IsColumnTable = 7; + static constexpr ui64 SysParam_TraversalStartKey = 2; + static constexpr ui64 SysParam_TraversalTableOwnerId = 3; + static constexpr ui64 SysParam_TraversalTableLocalPathId = 4; + static constexpr ui64 SysParam_TraversalStartTime = 5; + static constexpr ui64 SysParam_LastForceTraversalOperationId = 6; + static constexpr ui64 SysParam_TraversalIsColumnTable = 7; static constexpr ui64 SysParam_GlobalTraversalRound = 8; }; diff --git a/ydb/core/statistics/aggregator/tx_analyze_table.cpp b/ydb/core/statistics/aggregator/tx_analyze_table.cpp new file mode 100644 index 00000000000..4fc965ac613 --- /dev/null +++ b/ydb/core/statistics/aggregator/tx_analyze_table.cpp @@ -0,0 +1,68 @@ +#include "aggregator_impl.h" + +#include <ydb/core/tx/datashard/datashard.h> + +namespace NKikimr::NStat { + +struct TStatisticsAggregator::TTxAnalyzeTable : public TTxBase { + TPathId PathId; + TActorId ReplyToActorId; + ui64 OperationId = 0; + + TTxAnalyzeTable(TSelf* self, const TPathId& pathId, TActorId replyToActorId) + : TTxBase(self) + , PathId(pathId) + , ReplyToActorId(replyToActorId) + {} + + TTxType GetTxType() const override { return TXTYPE_ANALYZE_TABLE; } + + bool Execute(TTransactionContext& txc, const TActorContext&) override { + SA_LOG_D("[" << Self->TabletID() << "] TTxAnalyzeTable::Execute"); + + if (!Self->EnableColumnStatistics) { + return true; + } + + auto itOp = Self->ForceTraversalsByPathId.find(PathId); + if (itOp != Self->ForceTraversalsByPathId.end()) { + itOp->second.ReplyToActorIds.insert(ReplyToActorId); + OperationId = itOp->second.OperationId; + return true; + } + + NIceDb::TNiceDb db(txc.DB); + + TForceTraversal& operation = Self->ForceTraversalsByPathId[PathId]; + operation.PathId = PathId; + operation.OperationId = ++Self->LastForceTraversalOperationId; + operation.ReplyToActorIds.insert(ReplyToActorId); + Self->ForceTraversals.PushBack(&operation); + + Self->PersistLastForceTraversalOperationId(db); + + db.Table<Schema::ForceTraversals>().Key(operation.OperationId).Update( + NIceDb::TUpdate<Schema::ForceTraversals::OwnerId>(PathId.OwnerId), + NIceDb::TUpdate<Schema::ForceTraversals::LocalPathId>(PathId.LocalPathId)); + + OperationId = operation.OperationId; + + return true; + } + + void Complete(const TActorContext& ctx) override { + SA_LOG_D("[" << Self->TabletID() << "] TTxAnalyzeTable::Complete"); + } +}; + +void TStatisticsAggregator::Handle(TEvStatistics::TEvAnalyze::TPtr& ev) { + const auto& record = ev->Get()->Record; + + // TODO: replace by queue + for (const auto& table : record.GetTables()) { + Execute(new TTxAnalyzeTable(this, PathIdFromPathId(table.GetPathId()), ev->Sender), TActivationContext::AsActorContext()); + } + +} + +} // NKikimr::NStat diff --git a/ydb/core/statistics/aggregator/tx_statistics_scan_response.cpp b/ydb/core/statistics/aggregator/tx_datashard_scan_response.cpp index f32a78450ef..79ea60b5bbc 100644 --- a/ydb/core/statistics/aggregator/tx_statistics_scan_response.cpp +++ b/ydb/core/statistics/aggregator/tx_datashard_scan_response.cpp @@ -4,11 +4,11 @@ namespace NKikimr::NStat { -struct TStatisticsAggregator::TTxStatisticsScanResponse : public TTxBase { +struct TStatisticsAggregator::TTxDatashardScanResponse : public TTxBase { NKikimrStat::TEvStatisticsResponse Record; bool IsCorrectShardId = false; - TTxStatisticsScanResponse(TSelf* self, NKikimrStat::TEvStatisticsResponse&& record) + TTxDatashardScanResponse(TSelf* self, NKikimrStat::TEvStatisticsResponse&& record) : TTxBase(self) , Record(std::move(record)) {} @@ -16,17 +16,17 @@ struct TStatisticsAggregator::TTxStatisticsScanResponse : public TTxBase { TTxType GetTxType() const override { return TXTYPE_SCAN_RESPONSE; } bool Execute(TTransactionContext& txc, const TActorContext&) override { - SA_LOG_D("[" << Self->TabletID() << "] TTxStatisticsScanResponse::Execute"); + SA_LOG_D("[" << Self->TabletID() << "] TTxDatashardScanResponse::Execute"); NIceDb::TNiceDb db(txc.DB); // TODO: handle scan errors - if (Self->ShardRanges.empty()) { + if (Self->DatashardRanges.empty()) { return true; } - auto& range = Self->ShardRanges.front(); + auto& range = Self->DatashardRanges.front(); auto replyShardId = Record.GetShardTabletId(); if (replyShardId != range.DataShardId) { @@ -53,31 +53,31 @@ struct TStatisticsAggregator::TTxStatisticsScanResponse : public TTxBase { *current += *sketch; auto currentStr = TString(current->AsStringBuf()); - db.Table<Schema::Statistics>().Key(tag).Update( - NIceDb::TUpdate<Schema::Statistics::CountMinSketch>(currentStr)); + db.Table<Schema::ColumnStatistics>().Key(tag).Update( + NIceDb::TUpdate<Schema::ColumnStatistics::CountMinSketch>(currentStr)); } } } - Self->StartKey = range.EndKey; + Self->TraversalStartKey = range.EndKey; Self->PersistStartKey(db); return true; } void Complete(const TActorContext&) override { - SA_LOG_D("[" << Self->TabletID() << "] TTxStatisticsScanResponse::Complete"); + SA_LOG_D("[" << Self->TabletID() << "] TTxDatashardScanResponse::Complete"); - if (IsCorrectShardId && !Self->ShardRanges.empty()) { - Self->ShardRanges.pop_front(); - Self->NextRange(); + if (IsCorrectShardId && !Self->DatashardRanges.empty()) { + Self->DatashardRanges.pop_front(); + Self->ScanNextDatashardRange(); } } }; void TStatisticsAggregator::Handle(NStat::TEvStatistics::TEvStatisticsResponse::TPtr& ev) { auto& record = ev->Get()->Record; - Execute(new TTxStatisticsScanResponse(this, std::move(record)), + Execute(new TTxDatashardScanResponse(this, std::move(record)), TActivationContext::AsActorContext()); } diff --git a/ydb/core/statistics/aggregator/tx_delete_query_response.cpp b/ydb/core/statistics/aggregator/tx_delete_query_response.cpp index 921fb6d2eb0..31fff8be20f 100644 --- a/ydb/core/statistics/aggregator/tx_delete_query_response.cpp +++ b/ydb/core/statistics/aggregator/tx_delete_query_response.cpp @@ -19,7 +19,7 @@ struct TStatisticsAggregator::TTxDeleteQueryResponse : public TTxBase { ReplyToActorIds.swap(Self->ReplyToActorIds); NIceDb::TNiceDb db(txc.DB); - Self->FinishScan(db); + Self->FinishTraversal(db); return true; } diff --git a/ydb/core/statistics/aggregator/tx_init.cpp b/ydb/core/statistics/aggregator/tx_init.cpp index 816c1168888..1f0df4f9e52 100644 --- a/ydb/core/statistics/aggregator/tx_init.cpp +++ b/ydb/core/statistics/aggregator/tx_init.cpp @@ -19,16 +19,16 @@ struct TStatisticsAggregator::TTxInit : public TTxBase { { // precharge auto sysParamsRowset = db.Table<Schema::SysParams>().Range().Select(); - auto baseStatsRowset = db.Table<Schema::BaseStats>().Range().Select(); - auto statisticsRowset = db.Table<Schema::Statistics>().Range().Select(); - auto scanTablesRowset = db.Table<Schema::ScanTables>().Range().Select(); - auto scanOperationsRowset = db.Table<Schema::ScanOperations>().Range().Select(); + auto baseStatisticsRowset = db.Table<Schema::BaseStatistics>().Range().Select(); + auto statisticsRowset = db.Table<Schema::ColumnStatistics>().Range().Select(); + auto scheduleTraversalRowset = db.Table<Schema::ScheduleTraversals>().Range().Select(); + auto forceTraversalRowset = db.Table<Schema::ForceTraversals>().Range().Select(); if (!sysParamsRowset.IsReady() || - !baseStatsRowset.IsReady() || + !baseStatisticsRowset.IsReady() || !statisticsRowset.IsReady() || - !scanTablesRowset.IsReady() || - !scanOperationsRowset.IsReady()) + !scheduleTraversalRowset.IsReady() || + !forceTraversalRowset.IsReady()) { return false; } @@ -48,41 +48,41 @@ struct TStatisticsAggregator::TTxInit : public TTxBase { switch (id) { case Schema::SysParam_Database: Self->Database = value; - SA_LOG_D("[" << Self->TabletID() << "] Loading database: " << Self->Database); + SA_LOG_D("[" << Self->TabletID() << "] Loaded database: " << Self->Database); break; - case Schema::SysParam_StartKey: - Self->StartKey = TSerializedCellVec(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading start key"); + case Schema::SysParam_TraversalStartKey: + Self->TraversalStartKey = TSerializedCellVec(value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded traversal start key"); break; - case Schema::SysParam_ScanTableOwnerId: - Self->ScanTableId.PathId.OwnerId = FromString<ui64>(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading scan table owner id: " - << Self->ScanTableId.PathId.OwnerId); + case Schema::SysParam_TraversalTableOwnerId: + Self->TraversalTableId.PathId.OwnerId = FromString<ui64>(value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded traversal table owner id: " + << Self->TraversalTableId.PathId.OwnerId); break; - case Schema::SysParam_ScanTableLocalPathId: - Self->ScanTableId.PathId.LocalPathId = FromString<ui64>(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading scan table local path id: " - << Self->ScanTableId.PathId.LocalPathId); + case Schema::SysParam_TraversalTableLocalPathId: + Self->TraversalTableId.PathId.LocalPathId = FromString<ui64>(value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded traversal table local path id: " + << Self->TraversalTableId.PathId.LocalPathId); break; - case Schema::SysParam_ScanStartTime: { + case Schema::SysParam_TraversalStartTime: { auto us = FromString<ui64>(value); - Self->ScanStartTime = TInstant::MicroSeconds(us); - SA_LOG_D("[" << Self->TabletID() << "] Loading scan start time: " << us); + Self->TraversalStartTime = TInstant::MicroSeconds(us); + SA_LOG_D("[" << Self->TabletID() << "] Loaded traversal start time: " << us); break; } - case Schema::SysParam_LastScanOperationId: { - Self->LastScanOperationId = FromString<ui64>(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading last scan operation id: " << value); + case Schema::SysParam_LastForceTraversalOperationId: { + Self->LastForceTraversalOperationId = FromString<ui64>(value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded last traversal operation id: " << value); break; } - case Schema::SysParam_IsColumnTable: { - Self->IsColumnTable = FromString<bool>(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading IsColumnTable: " << value); + case Schema::SysParam_TraversalIsColumnTable: { + Self->TraversalIsColumnTable = FromString<bool>(value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded traversal IsColumnTable: " << value); break; } case Schema::SysParam_GlobalTraversalRound: { Self->GlobalTraversalRound = FromString<ui64>(value); - SA_LOG_D("[" << Self->TabletID() << "] Loading global traversal round: " << value); + SA_LOG_D("[" << Self->TabletID() << "] Loaded global traversal round: " << value); break; } default: @@ -95,42 +95,42 @@ struct TStatisticsAggregator::TTxInit : public TTxBase { } } - // BaseStats + // BaseStatistics { - Self->BaseStats.clear(); + Self->BaseStatistics.clear(); - auto rowset = db.Table<Schema::BaseStats>().Range().Select(); + auto rowset = db.Table<Schema::BaseStatistics>().Range().Select(); if (!rowset.IsReady()) { return false; } while (!rowset.EndOfSet()) { - ui64 schemeShardId = rowset.GetValue<Schema::BaseStats::SchemeShardId>(); - TString stats = rowset.GetValue<Schema::BaseStats::Stats>(); + ui64 schemeShardId = rowset.GetValue<Schema::BaseStatistics::SchemeShardId>(); + TString stats = rowset.GetValue<Schema::BaseStatistics::Stats>(); - Self->BaseStats[schemeShardId] = stats; + Self->BaseStatistics[schemeShardId] = stats; if (!rowset.Next()) { return false; } } - SA_LOG_D("[" << Self->TabletID() << "] Loading base stats: " - << "schemeshard count# " << Self->BaseStats.size()); + SA_LOG_D("[" << Self->TabletID() << "] Loaded BaseStatistics: " + << "schemeshard count# " << Self->BaseStatistics.size()); } - // Statistics + // ColumnStatistics { Self->CountMinSketches.clear(); - auto rowset = db.Table<Schema::Statistics>().Range().Select(); + auto rowset = db.Table<Schema::ColumnStatistics>().Range().Select(); if (!rowset.IsReady()) { return false; } while (!rowset.EndOfSet()) { - ui32 columnTag = rowset.GetValue<Schema::Statistics::ColumnTag>(); - TString sketch = rowset.GetValue<Schema::Statistics::CountMinSketch>(); + ui32 columnTag = rowset.GetValue<Schema::ColumnStatistics::ColumnTag>(); + TString sketch = rowset.GetValue<Schema::ColumnStatistics::CountMinSketch>(); Self->CountMinSketches[columnTag].reset( TCountMinSketch::FromString(sketch.data(), sketch.size())); @@ -140,78 +140,78 @@ struct TStatisticsAggregator::TTxInit : public TTxBase { } } - SA_LOG_D("[" << Self->TabletID() << "] Loading statistics: " + SA_LOG_D("[" << Self->TabletID() << "] Loaded ColumnStatistics: " << "column count# " << Self->CountMinSketches.size()); } - // ScanTables + // ScheduleTraversals { - Self->ScanTablesByTime.Clear(); - Self->ScanTablesBySchemeShard.clear(); - Self->ScanTables.clear(); + Self->ScheduleTraversalsByTime.Clear(); + Self->ScheduleTraversalsBySchemeShard.clear(); + Self->ScheduleTraversals.clear(); - auto rowset = db.Table<Schema::ScanTables>().Range().Select(); + auto rowset = db.Table<Schema::ScheduleTraversals>().Range().Select(); if (!rowset.IsReady()) { return false; } while (!rowset.EndOfSet()) { - ui64 ownerId = rowset.GetValue<Schema::ScanTables::OwnerId>(); - ui64 localPathId = rowset.GetValue<Schema::ScanTables::LocalPathId>(); - ui64 lastUpdateTime = rowset.GetValue<Schema::ScanTables::LastUpdateTime>(); - ui64 schemeShardId = rowset.GetValue<Schema::ScanTables::SchemeShardId>(); - bool isColumnTable = rowset.GetValue<Schema::ScanTables::IsColumnTable>(); + ui64 ownerId = rowset.GetValue<Schema::ScheduleTraversals::OwnerId>(); + ui64 localPathId = rowset.GetValue<Schema::ScheduleTraversals::LocalPathId>(); + ui64 lastUpdateTime = rowset.GetValue<Schema::ScheduleTraversals::LastUpdateTime>(); + ui64 schemeShardId = rowset.GetValue<Schema::ScheduleTraversals::SchemeShardId>(); + bool isColumnTable = rowset.GetValue<Schema::ScheduleTraversals::IsColumnTable>(); auto pathId = TPathId(ownerId, localPathId); - TScanTable scanTable; - scanTable.PathId = pathId; - scanTable.SchemeShardId = schemeShardId; - scanTable.LastUpdateTime = TInstant::MicroSeconds(lastUpdateTime); - scanTable.IsColumnTable = isColumnTable; + TScheduleTraversal scheduleTraversal; + scheduleTraversal.PathId = pathId; + scheduleTraversal.SchemeShardId = schemeShardId; + scheduleTraversal.LastUpdateTime = TInstant::MicroSeconds(lastUpdateTime); + scheduleTraversal.IsColumnTable = isColumnTable; - auto [it, _] = Self->ScanTables.emplace(pathId, scanTable); - Self->ScanTablesByTime.Add(&it->second); - Self->ScanTablesBySchemeShard[schemeShardId].insert(pathId); + auto [it, _] = Self->ScheduleTraversals.emplace(pathId, scheduleTraversal); + Self->ScheduleTraversalsByTime.Add(&it->second); + Self->ScheduleTraversalsBySchemeShard[schemeShardId].insert(pathId); if (!rowset.Next()) { return false; } } - SA_LOG_D("[" << Self->TabletID() << "] Loading scan tables: " - << "table count# " << Self->ScanTables.size()); + SA_LOG_D("[" << Self->TabletID() << "] Loaded ScheduleTraversals: " + << "table count# " << Self->ScheduleTraversals.size()); } - // ScanOperations + // ForceTraversals { - Self->ScanOperations.Clear(); - Self->ScanOperationsByPathId.clear(); + Self->ForceTraversals.Clear(); + Self->ForceTraversalsByPathId.clear(); - auto rowset = db.Table<Schema::ScanOperations>().Range().Select(); + auto rowset = db.Table<Schema::ForceTraversals>().Range().Select(); if (!rowset.IsReady()) { return false; } while (!rowset.EndOfSet()) { - ui64 operationId = rowset.GetValue<Schema::ScanOperations::OperationId>(); - ui64 ownerId = rowset.GetValue<Schema::ScanOperations::OwnerId>(); - ui64 localPathId = rowset.GetValue<Schema::ScanOperations::LocalPathId>(); + ui64 operationId = rowset.GetValue<Schema::ForceTraversals::OperationId>(); + ui64 ownerId = rowset.GetValue<Schema::ForceTraversals::OwnerId>(); + ui64 localPathId = rowset.GetValue<Schema::ForceTraversals::LocalPathId>(); auto pathId = TPathId(ownerId, localPathId); - TScanOperation& operation = Self->ScanOperationsByPathId[pathId]; + TForceTraversal& operation = Self->ForceTraversalsByPathId[pathId]; operation.PathId = pathId; operation.OperationId = operationId; - Self->ScanOperations.PushBack(&operation); + Self->ForceTraversals.PushBack(&operation); if (!rowset.Next()) { return false; } } - SA_LOG_D("[" << Self->TabletID() << "] Loading scan operations: " - << "table count# " << Self->ScanOperationsByPathId.size()); + SA_LOG_D("[" << Self->TabletID() << "] Loaded ForceTraversals: " + << "table count# " << Self->ForceTraversalsByPathId.size()); } return true; @@ -227,11 +227,11 @@ struct TStatisticsAggregator::TTxInit : public TTxBase { Self->SubscribeForConfigChanges(ctx); Self->Schedule(Self->PropagateInterval, new TEvPrivate::TEvPropagate()); - Self->Schedule(Self->ScheduleScanIntervalTime, new TEvPrivate::TEvScheduleScan()); + Self->Schedule(Self->TraversalPeriod, new TEvPrivate::TEvScheduleTraversal()); Self->InitializeStatisticsTable(); - if (Self->ScanTableId.PathId) { + if (Self->TraversalTableId.PathId) { Self->Navigate(); } diff --git a/ydb/core/statistics/aggregator/tx_init_schema.cpp b/ydb/core/statistics/aggregator/tx_init_schema.cpp index 17da2827e3c..35278624558 100644 --- a/ydb/core/statistics/aggregator/tx_init_schema.cpp +++ b/ydb/core/statistics/aggregator/tx_init_schema.cpp @@ -15,9 +15,9 @@ struct TStatisticsAggregator::TTxInitSchema : public TTxBase { NIceDb::TNiceDb(txc.DB).Materialize<Schema>(); static constexpr NIceDb::TTableId bigTableIds[] = { - Schema::BaseStats::TableId, - Schema::Statistics::TableId, - Schema::ScanTables::TableId + Schema::BaseStatistics::TableId, + Schema::ColumnStatistics::TableId, + Schema::ScheduleTraversals::TableId }; for (auto id : bigTableIds) { diff --git a/ydb/core/statistics/aggregator/tx_navigate.cpp b/ydb/core/statistics/aggregator/tx_navigate.cpp index bbc2ebdf6da..35534d042b0 100644 --- a/ydb/core/statistics/aggregator/tx_navigate.cpp +++ b/ydb/core/statistics/aggregator/tx_navigate.cpp @@ -29,7 +29,7 @@ struct TStatisticsAggregator::TTxNavigate : public TTxBase { if (entry.Status == NSchemeCache::TSchemeCacheNavigate::EStatus::PathErrorUnknown) { Self->DeleteStatisticsFromTable(); } else { - Self->FinishScan(db); + Self->FinishTraversal(db); } return true; } @@ -52,13 +52,13 @@ struct TStatisticsAggregator::TTxNavigate : public TTxBase { Self->KeyColumnTypes[col.second.KeyOrder] = col.second.PType; } - if (Self->StartKey.GetCells().empty()) { + if (Self->TraversalStartKey.GetCells().empty()) { TVector<TCell> minusInf(Self->KeyColumnTypes.size()); - Self->StartKey = TSerializedCellVec(minusInf); + Self->TraversalStartKey = TSerializedCellVec(minusInf); Self->PersistStartKey(db); } - if (Self->IsColumnTable) { + if (Self->TraversalIsColumnTable) { // TODO: serverless case if (entry.DomainInfo->Params.HasHive()) { Self->HiveId = entry.DomainInfo->Params.GetHive(); diff --git a/ydb/core/statistics/aggregator/tx_resolve.cpp b/ydb/core/statistics/aggregator/tx_resolve.cpp index b0c6b2bbffb..2c94c5e6d33 100644 --- a/ydb/core/statistics/aggregator/tx_resolve.cpp +++ b/ydb/core/statistics/aggregator/tx_resolve.cpp @@ -30,30 +30,30 @@ struct TStatisticsAggregator::TTxResolve : public TTxBase { if (entry.Status == NSchemeCache::TSchemeCacheRequest::EStatus::PathErrorNotExist) { Self->DeleteStatisticsFromTable(); } else { - Self->FinishScan(db); + Self->FinishTraversal(db); } return true; } auto& partitioning = entry.KeyDescription->GetPartitions(); - if (Self->IsColumnTable) { + if (Self->TraversalIsColumnTable) { Self->TabletsForReqDistribution.clear(); } else { - Self->ShardRanges.clear(); + Self->DatashardRanges.clear(); } for (auto& part : partitioning) { if (!part.Range) { continue; } - if (Self->IsColumnTable) { + if (Self->TraversalIsColumnTable) { Self->TabletsForReqDistribution.insert(part.ShardId); } else { TRange range; range.EndKey = part.Range->EndKeyPrefix; range.DataShardId = part.ShardId; - Self->ShardRanges.push_back(range); + Self->DatashardRanges.push_back(range); } } @@ -67,10 +67,10 @@ struct TStatisticsAggregator::TTxResolve : public TTxBase { return; } - if (Self->IsColumnTable) { + if (Self->TraversalIsColumnTable) { ctx.Send(Self->SelfId(), new TEvPrivate::TEvRequestDistribution); } else { - Self->NextRange(); + Self->ScanNextDatashardRange(); } } }; diff --git a/ydb/core/statistics/aggregator/tx_response_tablet_distribution.cpp b/ydb/core/statistics/aggregator/tx_response_tablet_distribution.cpp index c54d5d4d951..b97c44eaa6e 100644 --- a/ydb/core/statistics/aggregator/tx_response_tablet_distribution.cpp +++ b/ydb/core/statistics/aggregator/tx_response_tablet_distribution.cpp @@ -33,7 +33,7 @@ struct TStatisticsAggregator::TTxResponseTabletDistribution : public TTxBase { Request = std::make_unique<TEvStatistics::TEvAggregateStatistics>(); auto& outRecord = Request->Record; - PathIdFromPathId(Self->ScanTableId.PathId, outRecord.MutablePathId()); + PathIdFromPathId(Self->TraversalTableId.PathId, outRecord.MutablePathId()); bool hasTablets = false; for (auto& inNode : Record.GetNodes()) { @@ -62,7 +62,7 @@ struct TStatisticsAggregator::TTxResponseTabletDistribution : public TTxBase { } if (!hasTablets) { - Self->FinishScan(db); + Self->FinishTraversal(db); return true; } diff --git a/ydb/core/statistics/aggregator/tx_save_query_response.cpp b/ydb/core/statistics/aggregator/tx_save_query_response.cpp index 76bf2d1d94a..19804f4f466 100644 --- a/ydb/core/statistics/aggregator/tx_save_query_response.cpp +++ b/ydb/core/statistics/aggregator/tx_save_query_response.cpp @@ -19,7 +19,7 @@ struct TStatisticsAggregator::TTxSaveQueryResponse : public TTxBase { ReplyToActorIds.swap(Self->ReplyToActorIds); NIceDb::TNiceDb db(txc.DB); - Self->FinishScan(db); + Self->FinishTraversal(db); return true; } diff --git a/ydb/core/statistics/aggregator/tx_scan_table.cpp b/ydb/core/statistics/aggregator/tx_scan_table.cpp deleted file mode 100644 index c1323bb3b22..00000000000 --- a/ydb/core/statistics/aggregator/tx_scan_table.cpp +++ /dev/null @@ -1,68 +0,0 @@ -#include "aggregator_impl.h" - -#include <ydb/core/tx/datashard/datashard.h> - -namespace NKikimr::NStat { - -struct TStatisticsAggregator::TTxScanTable : public TTxBase { - TPathId PathId; - TActorId ReplyToActorId; - ui64 OperationId = 0; - - TTxScanTable(TSelf* self, const TPathId& pathId, TActorId replyToActorId) - : TTxBase(self) - , PathId(pathId) - , ReplyToActorId(replyToActorId) - {} - - TTxType GetTxType() const override { return TXTYPE_SCAN_TABLE; } - - bool Execute(TTransactionContext& txc, const TActorContext&) override { - SA_LOG_D("[" << Self->TabletID() << "] TTxScanTable::Execute"); - - if (!Self->EnableColumnStatistics) { - return true; - } - - auto itOp = Self->ScanOperationsByPathId.find(PathId); - if (itOp != Self->ScanOperationsByPathId.end()) { - itOp->second.ReplyToActorIds.insert(ReplyToActorId); - OperationId = itOp->second.OperationId; - return true; - } - - NIceDb::TNiceDb db(txc.DB); - - TScanOperation& operation = Self->ScanOperationsByPathId[PathId]; - operation.PathId = PathId; - operation.OperationId = ++Self->LastScanOperationId; - operation.ReplyToActorIds.insert(ReplyToActorId); - Self->ScanOperations.PushBack(&operation); - - Self->PersistLastScanOperationId(db); - - db.Table<Schema::ScanOperations>().Key(operation.OperationId).Update( - NIceDb::TUpdate<Schema::ScanOperations::OwnerId>(PathId.OwnerId), - NIceDb::TUpdate<Schema::ScanOperations::LocalPathId>(PathId.LocalPathId)); - - OperationId = operation.OperationId; - - return true; - } - - void Complete(const TActorContext& ctx) override { - SA_LOG_D("[" << Self->TabletID() << "] TTxScanTable::Complete"); - } -}; - -void TStatisticsAggregator::Handle(TEvStatistics::TEvAnalyze::TPtr& ev) { - const auto& record = ev->Get()->Record; - - // TODO: replace by queue - for (const auto& table : record.GetTables()) { - Execute(new TTxScanTable(this, PathIdFromPathId(table.GetPathId()), ev->Sender), TActivationContext::AsActorContext()); - } - -} - -} // NKikimr::NStat diff --git a/ydb/core/statistics/aggregator/tx_schedule_scan.cpp b/ydb/core/statistics/aggregator/tx_schedule_scan.cpp deleted file mode 100644 index 9c3e1f6adae..00000000000 --- a/ydb/core/statistics/aggregator/tx_schedule_scan.cpp +++ /dev/null @@ -1,41 +0,0 @@ -#include "aggregator_impl.h" - -#include <ydb/core/tx/datashard/datashard.h> - -namespace NKikimr::NStat { - -struct TStatisticsAggregator::TTxScheduleScan : public TTxBase { - TTxScheduleScan(TSelf* self) - : TTxBase(self) - {} - - TTxType GetTxType() const override { return TXTYPE_SCHEDULE_SCAN; } - - bool Execute(TTransactionContext& txc, const TActorContext&) override { - SA_LOG_T("[" << Self->TabletID() << "] TTxScheduleScan::Execute"); - - Self->Schedule(Self->ScheduleScanIntervalTime, new TEvPrivate::TEvScheduleScan()); - - if (!Self->EnableColumnStatistics) { - return true; - } - - if (Self->ScanTableId.PathId) { - return true; // scan is in progress - } - - NIceDb::TNiceDb db(txc.DB); - Self->ScheduleNextScan(db); - return true; - } - - void Complete(const TActorContext&) override { - SA_LOG_T("[" << Self->TabletID() << "] TTxScheduleScan::Complete"); - } -}; - -void TStatisticsAggregator::Handle(TEvPrivate::TEvScheduleScan::TPtr&) { - Execute(new TTxScheduleScan(this), TActivationContext::AsActorContext()); -} - -} // NKikimr::NStat diff --git a/ydb/core/statistics/aggregator/tx_schedule_traversal.cpp b/ydb/core/statistics/aggregator/tx_schedule_traversal.cpp new file mode 100644 index 00000000000..0eb09bbd67a --- /dev/null +++ b/ydb/core/statistics/aggregator/tx_schedule_traversal.cpp @@ -0,0 +1,41 @@ +#include "aggregator_impl.h" + +#include <ydb/core/tx/datashard/datashard.h> + +namespace NKikimr::NStat { + +struct TStatisticsAggregator::TTxScheduleTrasersal : public TTxBase { + TTxScheduleTrasersal(TSelf* self) + : TTxBase(self) + {} + + TTxType GetTxType() const override { return TXTYPE_SCHEDULE_TRAVERSAL; } + + bool Execute(TTransactionContext& txc, const TActorContext&) override { + SA_LOG_T("[" << Self->TabletID() << "] TTxScheduleTrasersal::Execute"); + + Self->Schedule(Self->TraversalPeriod, new TEvPrivate::TEvScheduleTraversal()); + + if (!Self->EnableColumnStatistics) { + return true; + } + + if (Self->TraversalTableId.PathId) { + return true; // traverse is in progress + } + + NIceDb::TNiceDb db(txc.DB); + Self->ScheduleNextTraversal(db); + return true; + } + + void Complete(const TActorContext&) override { + SA_LOG_T("[" << Self->TabletID() << "] TTxScheduleTrasersal::Complete"); + } +}; + +void TStatisticsAggregator::Handle(TEvPrivate::TEvScheduleTraversal::TPtr&) { + Execute(new TTxScheduleTrasersal(this), TActivationContext::AsActorContext()); +} + +} // NKikimr::NStat diff --git a/ydb/core/statistics/aggregator/tx_schemeshard_stats.cpp b/ydb/core/statistics/aggregator/tx_schemeshard_stats.cpp index 6dc87267e57..61efd9a2c74 100644 --- a/ydb/core/statistics/aggregator/tx_schemeshard_stats.cpp +++ b/ydb/core/statistics/aggregator/tx_schemeshard_stats.cpp @@ -21,10 +21,10 @@ struct TStatisticsAggregator::TTxSchemeShardStats : public TTxBase { << ", stats size# " << stats.size()); NIceDb::TNiceDb db(txc.DB); - db.Table<Schema::BaseStats>().Key(schemeShardId).Update( - NIceDb::TUpdate<Schema::BaseStats::Stats>(stats)); + db.Table<Schema::BaseStatistics>().Key(schemeShardId).Update( + NIceDb::TUpdate<Schema::BaseStatistics::Stats>(stats)); - Self->BaseStats[schemeShardId] = stats; + Self->BaseStatistics[schemeShardId] = stats; if (!Self->EnableColumnStatistics) { return true; @@ -33,39 +33,39 @@ struct TStatisticsAggregator::TTxSchemeShardStats : public TTxBase { NKikimrStat::TSchemeShardStats statRecord; Y_PROTOBUF_SUPPRESS_NODISCARD statRecord.ParseFromString(stats); - auto& oldPathIds = Self->ScanTablesBySchemeShard[schemeShardId]; + auto& oldPathIds = Self->ScheduleTraversalsBySchemeShard[schemeShardId]; std::unordered_set<TPathId> newPathIds; for (auto& entry : statRecord.GetEntries()) { auto pathId = PathIdFromPathId(entry.GetPathId()); newPathIds.insert(pathId); if (oldPathIds.find(pathId) == oldPathIds.end()) { - TStatisticsAggregator::TScanTable scanTable; - scanTable.PathId = pathId; - scanTable.SchemeShardId = schemeShardId; - scanTable.LastUpdateTime = TInstant::MicroSeconds(0); - scanTable.IsColumnTable = entry.GetIsColumnTable(); - auto [it, _] = Self->ScanTables.emplace(pathId, scanTable); - if (!Self->ScanTablesByTime.Has(&it->second)) { - Self->ScanTablesByTime.Add(&it->second); + TStatisticsAggregator::TScheduleTraversal traversalTable; + traversalTable.PathId = pathId; + traversalTable.SchemeShardId = schemeShardId; + traversalTable.LastUpdateTime = TInstant::MicroSeconds(0); + traversalTable.IsColumnTable = entry.GetIsColumnTable(); + auto [it, _] = Self->ScheduleTraversals.emplace(pathId, traversalTable); + if (!Self->ScheduleTraversalsByTime.Has(&it->second)) { + Self->ScheduleTraversalsByTime.Add(&it->second); } - db.Table<Schema::ScanTables>().Key(pathId.OwnerId, pathId.LocalPathId).Update( - NIceDb::TUpdate<Schema::ScanTables::SchemeShardId>(schemeShardId), - NIceDb::TUpdate<Schema::ScanTables::LastUpdateTime>(0), - NIceDb::TUpdate<Schema::ScanTables::IsColumnTable>(entry.GetIsColumnTable())); + db.Table<Schema::ScheduleTraversals>().Key(pathId.OwnerId, pathId.LocalPathId).Update( + NIceDb::TUpdate<Schema::ScheduleTraversals::SchemeShardId>(schemeShardId), + NIceDb::TUpdate<Schema::ScheduleTraversals::LastUpdateTime>(0), + NIceDb::TUpdate<Schema::ScheduleTraversals::IsColumnTable>(entry.GetIsColumnTable())); } } for (auto& pathId : oldPathIds) { if (newPathIds.find(pathId) == newPathIds.end()) { - auto it = Self->ScanTables.find(pathId); - if (it != Self->ScanTables.end()) { - if (Self->ScanTablesByTime.Has(&it->second)) { - Self->ScanTablesByTime.Remove(&it->second); + auto it = Self->ScheduleTraversals.find(pathId); + if (it != Self->ScheduleTraversals.end()) { + if (Self->ScheduleTraversalsByTime.Has(&it->second)) { + Self->ScheduleTraversalsByTime.Remove(&it->second); } - Self->ScanTables.erase(it); + Self->ScheduleTraversals.erase(it); } - db.Table<Schema::ScanTables>().Key(pathId.OwnerId, pathId.LocalPathId).Delete(); + db.Table<Schema::ScheduleTraversals>().Key(pathId.OwnerId, pathId.LocalPathId).Delete(); } } diff --git a/ydb/core/statistics/aggregator/ya.make b/ydb/core/statistics/aggregator/ya.make index bcde9537fb4..9917f39d777 100644 --- a/ydb/core/statistics/aggregator/ya.make +++ b/ydb/core/statistics/aggregator/ya.make @@ -9,7 +9,9 @@ SRCS( schema.cpp tx_ack_timeout.cpp tx_aggr_stat_response.cpp + tx_analyze_table.cpp tx_configure.cpp + tx_datashard_scan_response.cpp tx_delete_query_response.cpp tx_init.cpp tx_init_schema.cpp @@ -17,10 +19,8 @@ SRCS( tx_resolve.cpp tx_response_tablet_distribution.cpp tx_save_query_response.cpp - tx_scan_table.cpp - tx_schedule_scan.cpp + tx_schedule_traversal.cpp tx_schemeshard_stats.cpp - tx_statistics_scan_response.cpp ) PEERDIR( diff --git a/ydb/core/statistics/service/http_request.cpp b/ydb/core/statistics/service/http_request.cpp index 98107ea4994..6e7b28081f7 100644 --- a/ydb/core/statistics/service/http_request.cpp +++ b/ydb/core/statistics/service/http_request.cpp @@ -107,13 +107,13 @@ void THttpRequest::Handle(TEvStatistics::TEvAnalyzeStatusResponse::TPtr& ev) { HttpReply("Status is unspecified"); break; case NKikimrStat::TEvAnalyzeStatusResponse::STATUS_NO_OPERATION: - HttpReply("No scan operation"); + HttpReply("No analyze operation"); break; case NKikimrStat::TEvAnalyzeStatusResponse::STATUS_ENQUEUED: - HttpReply("Scan is enqueued"); + HttpReply("Analyze is enqueued"); break; case NKikimrStat::TEvAnalyzeStatusResponse::STATUS_IN_PROGRESS: - HttpReply("Scan is in progress"); + HttpReply("Analyze is in progress"); break; } } @@ -136,7 +136,7 @@ void THttpRequest::ResolveSuccess() { Send(MakePipePerNodeCacheID(false), new TEvPipeCache::TEvForward(analyze.release(), StatisticsAggregatorId, true)); - HttpReply("Scan sent"); + HttpReply("Analyze sent"); } else { auto getStatus = std::make_unique<TEvStatistics::TEvAnalyzeStatus>(); auto& record = getStatus->Record; |
