/
githubmirror
/
nbs
Обзор
Документация
Войти
/
githubmirror
/
nbs
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
cloud/filestore/libs/diagnostics/profile_log.cpp
246 строк
6 KB
Darya Frolova
issue-6751: add metric for compressed profile log frame size (#6752)
11 авг 2026, 14:02
Не верифицирован
11 авг 2026, 14:02
e1efdde
Код
Авторство
О чём код?
#include "profile_log.h" #include <cloud/storage/core/libs/common/scheduler.h> #include <cloud/storage/core/libs/common/timer.h> #include <library/cpp/eventlog/eventlog.h> #include <library/cpp/monlib/dynamic_counters/counters.h> #include <util/datetime/cputimer.h> #include <util/generic/array_ref.h> #include <util/generic/hash.h> #include <util/thread/lfstack.h> namespace NCloud::NFileStore { namespace { //////////////////////////////////////////////////////////////////////////////// using namespace NMonitoring; class TProfileLog final : public IProfileLog , public IEventLog::ISuccessCallback , public std::enable_shared_from_this<TProfileLog> { private: struct TCounters { TDynamicCounters::TCounterPtr Requests; TDynamicCounters::TCounterPtr Bytes; TDynamicCounters::TCounterPtr CompressedFrameBytes; TDynamicCounters::TCounterPtr Flushes; TDynamicCounters::TCounterPtr Discards; }; private: TEventLog EventLog; TProfileLogSettings Settings; ITimerPtr Timer; ISchedulerPtr Scheduler; TCounters Counters; bool CountersRegistered = false; TLockFreeStack<TRecord> Records; TAtomic ShouldStop = false; public: TProfileLog( TProfileLogSettings settings, ITimerPtr timer, ISchedulerPtr scheduler) : EventLog(settings.FilePath, NEvClass::Factory()->CurrentFormat()) , Settings(std::move(settings)) , Timer(std::move(timer)) , Scheduler(std::move(scheduler)) { EventLog.SetSuccessCallback(this); } ~TProfileLog() override; public: void Start() override; void Stop() override; void Write(TRecord record) override; void RegisterCounters(NMonitoring::TDynamicCounters& root) override; private: void OnWriteSuccess(const TBuffer& frameData) override; void ScheduleFlush(); void Flush(); }; TProfileLog::~TProfileLog() { Flush(); } void TProfileLog::Start() { ScheduleFlush(); } void TProfileLog::Stop() { AtomicSet(ShouldStop, 1); } void TProfileLog::Write(TRecord record) { if (CountersRegistered) { Counters.Requests->Inc(); Counters.Bytes->Add(static_cast<i64>(record.Request.ByteSizeLong())); } Records.Enqueue(std::move(record)); } void TProfileLog::RegisterCounters(NMonitoring::TDynamicCounters& root) { auto pg = root.GetSubgroup("counters", "profile_log"); Counters.Requests = pg->GetCounter("Count", true); Counters.Bytes = pg->GetCounter("RequestBytes", true); Counters.CompressedFrameBytes = pg->GetCounter("CompressedFrameBytes", true); Counters.Flushes = pg->GetCounter("FlushCount", true); Counters.Discards = pg->GetCounter("DiscardCount", true); CountersRegistered = true; } void TProfileLog::OnWriteSuccess(const TBuffer& frameData) { if (CountersRegistered) { Counters.CompressedFrameBytes->Add( static_cast<i64>(frameData.Size())); } } void TProfileLog::ScheduleFlush() { if (AtomicGet(ShouldStop)) { return; } Scheduler->Schedule( Timer->Now() + Settings.TimeThreshold, [weakPtr = weak_from_this()] { if (auto self = weakPtr.lock()) { self->Flush(); self->ScheduleFlush(); } } ); } void TProfileLog::Flush() { TVector<TRecord> records; Records.DequeueAllSingleConsumer(&records); ui64 discardedRequestCount = 0; if (Settings.MaxFlushRecords && records.size() > Settings.MaxFlushRecords) { discardedRequestCount = records.size() - Settings.MaxFlushRecords; records.resize(Settings.MaxFlushRecords); if (CountersRegistered) { Counters.Discards->Add(static_cast<i64>(discardedRequestCount)); } } if (CountersRegistered) { Counters.Flushes->Inc(); } auto recordsRef = MakeArrayRef(records); size_t frameRecordsOffset = 0; size_t maxFrameRecords = Settings.MaxFrameFlushRecords ? Settings.MaxFrameFlushRecords : records.size(); while (frameRecordsOffset < recordsRef.size()) { auto frameRecords = recordsRef.Slice( frameRecordsOffset, std::min(maxFrameRecords, recordsRef.size() - frameRecordsOffset)); TSelfFlushLogFrame logFrame(EventLog); THashMap<TString, TVector<ui32>> fsId2records; for (ui32 i = 0; i < frameRecords.size(); ++i) { fsId2records[frameRecords[i].FileSystemId].push_back(i); } for (auto& x: fsId2records) { NProto::TProfileLogRecord pb; if (discardedRequestCount) { pb.SetDiscardedRequestCount(discardedRequestCount); discardedRequestCount = 0; } pb.SetFileSystemId(x.first); for (const auto r: x.second) { auto& record = frameRecords[r]; *pb.AddRequests() = std::move(record.Request); } logFrame.LogEvent(pb); } frameRecordsOffset += frameRecords.size(); } EventLog.Flush(); } //////////////////////////////////////////////////////////////////////////////// class TProfileLogStub final : public IProfileLog { public: void Start() override { } void Stop() override { } void Write(TRecord record) override { Y_UNUSED(record); } void RegisterCounters(NMonitoring::TDynamicCounters& root) override { Y_UNUSED(root); } }; } // namespace //////////////////////////////////////////////////////////////////////////////// IProfileLogPtr CreateProfileLog( TProfileLogSettings settings, ITimerPtr timer, ISchedulerPtr scheduler) { return std::make_shared<TProfileLog>( std::move(settings), std::move(timer), std::move(scheduler) ); } IProfileLogPtr CreateProfileLogStub() { return std::make_shared<TProfileLogStub>(); } } // namespace NCloud::NFileStore