diff options
| author | kseleznyov <[email protected]> | 2026-07-15 17:59:21 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-15 17:59:21 +0300 |
| commit | 4cd9d470a6cc9034a46b8f19300220bc045a67f7 (patch) | |
| tree | a33727c3aaa81cb4a04d9e31db7a21909a41b395 | |
| parent | f3db5fee3486b7d66987cb0f05c03c1be7a91489 (diff) | |
[YDB_LOG] Migrate ydb/core/persqueue/public (#45809)
15 files changed, 340 insertions, 138 deletions
diff --git a/ydb/core/persqueue/public/cloud_events/actor.cpp b/ydb/core/persqueue/public/cloud_events/actor.cpp index 3d76983bf86..8f4f4b32bc2 100644 --- a/ydb/core/persqueue/public/cloud_events/actor.cpp +++ b/ydb/core/persqueue/public/cloud_events/actor.cpp @@ -19,6 +19,8 @@ #include <google/protobuf/util/json_util.h> #include <google/protobuf/util/time_util.h> +#define YDB_LOG_THIS_FILE_COMPONENT NKikimrServices::PERSQUEUE + namespace NKikimr::NPQ::NCloudEvents { using TCreateTopicEvent = yandex::cloud::events::ydb::topics::CreateTopic; @@ -325,8 +327,7 @@ template<typename TEvent> TString SerializeEventToProtobuf(const TEvent& ev) { TString data; if (!ev.SerializeToString(&data)) { - LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PERSQUEUE, - "SerializeToString failed"); + YDB_LOG_ERROR("SerializeToString failed"); return TString(); } return data; @@ -339,8 +340,7 @@ TString SerializeEventToJson(const TEvent& ev) { opts.preserve_proto_field_names = true; opts.always_print_primitive_fields = true; if (!google::protobuf::util::MessageToJsonString(ev, &data, opts).ok()) { - LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PERSQUEUE, - "MessageToJsonString failed"); + YDB_LOG_ERROR("MessageToJsonString failed"); return TString(); } return data; @@ -412,8 +412,7 @@ void TCloudEventsActor::Bootstrap() { void TCloudEventsActor::Handle(TCloudEvent::TPtr& ev) { if (!EventsWriter) { - LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PERSQUEUE, - "No events writer configured"); + YDB_LOG_ERROR("No events writer configured"); PassAway(); return; } @@ -421,15 +420,14 @@ void TCloudEventsActor::Handle(TCloudEvent::TPtr& ev) { try { TString data = BuildTopicCloudEvent(ev.Get()->Get()->Info); if (data.empty()) { - LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PERSQUEUE, - "Failed to build cloud event"); + YDB_LOG_ERROR("Failed to build cloud event"); PassAway(); return; } EventsWriter->Write(data); } catch (const std::exception& e) { - LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PERSQUEUE, - "Failed to write cloud event: " << e.what()); + YDB_LOG_ERROR("Failed to write cloud", + {"event", e.what()}); } PassAway(); diff --git a/ydb/core/persqueue/public/cluster_tracker/cluster_tracker.cpp b/ydb/core/persqueue/public/cluster_tracker/cluster_tracker.cpp index 59163c2bcd2..e482634a146 100644 --- a/ydb/core/persqueue/public/cluster_tracker/cluster_tracker.cpp +++ b/ydb/core/persqueue/public/cluster_tracker/cluster_tracker.cpp @@ -15,6 +15,8 @@ #include <ranges> #include <tuple> +#define YDB_LOG_THIS_FILE_COMPONENT NKikimrServices::PERSQUEUE_CLUSTER_TRACKER + namespace NKikimr::NPQ::NClusterTracker { inline auto& Ctx() { @@ -68,7 +70,8 @@ public: private: void AddSubscriber(const TActorId subscriberId) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "Subscribers.size: " << Subscribers.size() << " AddSubscriber"); + YDB_LOG_DEBUG_CTX(Ctx(), "AddSubscriber", + {"subscribersSize", Subscribers.size()}); Subscribers.insert(subscriberId); } @@ -81,7 +84,7 @@ private: } void HandleWhileWaiting(TEvClusterTracker::TEvSubscribe::TPtr& ev) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "AddSubscriber TEvSubscriber"); + YDB_LOG_DEBUG_CTX(Ctx(), "AddSubscriber TEvSubscriber"); // if (Cfg().GetTopicsAreFirstClassCitizen()) { // return SendClustersList(ev->Sender); @@ -94,7 +97,7 @@ private: } void HandleWhileWaiting(TEvClusterTracker::TEvGetClustersList::TPtr& ev) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "HandleWhileWaiting TEvGetClustersList"); + YDB_LOG_DEBUG_CTX(Ctx(), "HandleWhileWaiting TEvGetClustersList"); // if (Cfg().GetTopicsAreFirstClassCitizen()) { // return SendGetClustersListResponse(ev->Sender); @@ -134,7 +137,7 @@ private: } void SendClustersList(const TActorId& subscriberId) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "SendClustersList"); + YDB_LOG_DEBUG_CTX(Ctx(), "SendClustersList"); auto ev = MakeHolder<TEvClusterTracker::TEvClustersUpdate>(); @@ -145,7 +148,9 @@ private: } void HandleWhileWorking(TEvClusterTracker::TEvSubscribe::TPtr& ev) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "HandleWhileWorking TEvSubscribe Subscribers.size: " << Subscribers.size() << " ClustersList: " << (ClustersList == nullptr ? "null" : std::to_string(ClustersList->Clusters.size()))); + YDB_LOG_DEBUG_CTX(Ctx(), "HandleWhileWorking TEvSubscribe", + {"subscribersSize", Subscribers.size()}, + {"clustersList", (ClustersList == nullptr ? "null" : std::to_string(ClustersList->Clusters.size()))}); AddSubscriber(ev->Sender); @@ -156,7 +161,7 @@ private: } void HandleWhileWorking(TEvClusterTracker::TEvGetClustersList::TPtr& ev) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "HandleWhileWorking TEvGetClustersList"); + YDB_LOG_DEBUG_CTX(Ctx(), "HandleWhileWorking TEvGetClustersList"); if (ClustersList) { SendGetClustersListResponse(ev->Sender); @@ -166,7 +171,8 @@ private: } void SendGetClustersListResponse(const TActorId& senderId, bool success = true) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "SendGetClustersListResponse to " << senderId); + YDB_LOG_DEBUG_CTX(Ctx(), "SendGetClustersListResponse", + {"senderId", senderId}); auto ev = MakeHolder<TEvClusterTracker::TEvGetClustersListResponse>(); ev->Success = success; @@ -176,16 +182,19 @@ private: } void BroadcastClustersUpdate() { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "BroadcastClustersUpdate Subscribers.size: " << Subscribers.size()); + YDB_LOG_DEBUG_CTX(Ctx(), "BroadcastClustersUpdate", + {"subscribersSize", Subscribers.size()}); for (const auto& subscriberId : Subscribers) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "BroadcastClustersUpdate subscriberId: " << subscriberId); + YDB_LOG_DEBUG_CTX(Ctx(), "BroadcastClustersUpdate", + {"subscriberId", subscriberId}); SendClustersList(subscriberId); } } void ReplyAllGetClustersListRequests(bool success = true) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "ReplyAllGetClustersListRequests GetClustersListRequests.size: " << GetClustersListRequests.size()); + YDB_LOG_DEBUG_CTX(Ctx(), "ReplyAllGetClustersListRequests", + {"clustersListRequestsSize", GetClustersListRequests.size()}); for (const auto& requestId : GetClustersListRequests) { SendGetClustersListResponse(requestId, success); @@ -211,13 +220,13 @@ private: } void HandleWhileWorking(NKqp::TEvKqp::TEvQueryResponse::TPtr& ev) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "HandleWhileWorking TEvQueryResponse"); + YDB_LOG_DEBUG_CTX(Ctx(), "HandleWhileWorking TEvQueryResponse"); const auto& record = ev->Get()->Record; if (record.GetYdbStatus() == Ydb::StatusIds::SUCCESS) { NYdb::TResultSetParser parser(record.GetResponse().GetYdbResults(0)); if (parser.RowsCount()) { - LOG_DEBUG_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "HandleWhileWorking TEvQueryResponse UpdateClustersList"); + YDB_LOG_DEBUG_CTX(Ctx(), "HandleWhileWorking TEvQueryResponse UpdateClustersList"); UpdateClustersList(parser); AFL_ENSURE(ClustersList); @@ -232,7 +241,8 @@ private: } } - LOG_ERROR_S(Ctx(), NKikimrServices::PERSQUEUE_CLUSTER_TRACKER, "failed to list clusters: " << record); + YDB_LOG_ERROR_CTX(Ctx(), "Failed to list", + {"clusters", record}); ClustersList = nullptr; Schedule(TDuration::Seconds(Cfg().GetClustersUpdateTimeoutOnErrorSec()), new TEvents::TEvWakeup); diff --git a/ydb/core/persqueue/public/describer/describer.cpp b/ydb/core/persqueue/public/describer/describer.cpp index 360955d8b75..885d68d035c 100644 --- a/ydb/core/persqueue/public/describer/describer.cpp +++ b/ydb/core/persqueue/public/describer/describer.cpp @@ -2,11 +2,9 @@ #include <util/generic/algorithm.h> +#define YDB_LOG_THIS_FILE_COMPONENT NKikimrServices::PQ_DESCRIBER + #define LOG_PREFIX NActors::TlsActivationContext->AsActorContext().SelfID -#define LOG_E(stream) LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PQ_DESCRIBER, LOG_PREFIX << stream) -#define LOG_W(stream) LOG_WARN_S(*NActors::TlsActivationContext, NKikimrServices::PQ_DESCRIBER, LOG_PREFIX << stream) -#define LOG_I(stream) LOG_INFO_S(*NActors::TlsActivationContext, NKikimrServices::PQ_DESCRIBER, LOG_PREFIX << stream) -#define LOG_D(stream) LOG_DEBUG_S(*NActors::TlsActivationContext, NKikimrServices::PQ_DESCRIBER, LOG_PREFIX << stream) namespace NKikimr::NPQ::NDescriber { @@ -45,7 +43,10 @@ public: } void DoRequest(const std::unordered_set<TString>& topicPath) { - LOG_D("Create request [" << JoinRange(", ", topicPath.begin(), topicPath.end()) << "] with SyncVersion=" << RetryWithSyncVersion); + YDB_LOG_DEBUG("Create request with", + {"logPrefix", LOG_PREFIX}, + {"topicPaths", JoinRange(", ", topicPath.begin(), topicPath.end())}, + {"syncVersion", RetryWithSyncVersion}); auto schemeRequest = std::make_unique<TSchemeCacheNavigate>(1); schemeRequest->DatabaseName = DatabasePath; @@ -71,7 +72,8 @@ public: } void Handle(TEvTxProxySchemeCache::TEvNavigateKeySetResult::TPtr& ev) { - LOG_D("Handle TEvTxProxySchemeCache::TEvNavigateKeySetResult"); + YDB_LOG_DEBUG("Handle TEvTxProxySchemeCache::TEvNavigateKeySetResult", + {"logPrefix", LOG_PREFIX}); auto& result = ev->Get()->Request; std::unordered_set<TString> unknownPaths; @@ -98,12 +100,18 @@ public: case TSchemeCacheNavigate::EStatus::RootUnknown: { if (RetryWithSyncVersion) { if (entry.SecurityObject && !HasAccess(Settings, entry.SecurityObject)) { - LOG_D("Path '" << realPath << "' UNAUTHORIZED"); + YDB_LOG_DEBUG("Path UNAUTHORIZED", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); + Result[originalPath] = TTopicInfo{ .Status = EStatus::UNAUTHORIZED }; } else { - LOG_D("Path '" << realPath << "' not found"); + YDB_LOG_DEBUG("Path not found", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); + Result[originalPath] = TTopicInfo{ .Status = EStatus::NOT_FOUND }; @@ -114,7 +122,9 @@ public: break; } case TSchemeCacheNavigate::EStatus::AccessDenied: { - LOG_D("Path '" << realPath << "' ACCESS DENIED"); + YDB_LOG_DEBUG("Path ACCESS DENIED", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); Result[originalPath] = TTopicInfo{ .Status = EStatus::UNAUTHORIZED }; @@ -122,7 +132,10 @@ public: } case TSchemeCacheNavigate::EStatus::Ok: { if (entry.Kind == NSchemeCache::TSchemeCacheNavigate::KindCdcStream) { - LOG_D("Path '" << realPath << "' is a CDC"); + YDB_LOG_DEBUG("Path is CDC", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); + CDCPaths[TStringBuilder() << realPath << "/streamImpl"] = { .OriginalPath = originalPath, .CdcStreamName = entry.Self->Info.GetName() @@ -131,7 +144,9 @@ public: } else if (entry.Kind == TSchemeCacheNavigate::EKind::KindTopic) { if (!entry.PQGroupInfo || entry.PQGroupInfo->Description.GetBalancerTabletID() == 0) { if (RetryWithSyncVersion) { - LOG_D("Path '" << realPath << "' not found"); + YDB_LOG_DEBUG("Path not found", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); Result[originalPath] = TTopicInfo{ .Status = EStatus::NOT_FOUND }; @@ -140,13 +155,18 @@ public: } } else { if (!HasAccess(Settings, entry.SecurityObject)) { - LOG_D("Path '" << realPath << "' UNAUTHORIZED"); + YDB_LOG_DEBUG("Path UNAUTHORIZED", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); + Result[originalPath] = TTopicInfo{ .Status = entry.SecurityObject->CheckAccess(NACLib::EAccessRights::DescribeSchema, *Settings.UserToken) ? EStatus::UNAUTHORIZED_WITH_DESCRIBE_ACCESS : EStatus::UNAUTHORIZED }; } else { - LOG_D("Path '" << realPath << "' SUCCESS"); + YDB_LOG_DEBUG("Path SUCCESS", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); Result[originalPath] = TTopicInfo{ .Status = EStatus::SUCCESS, .RealPath = realPath, @@ -160,9 +180,14 @@ public: } } } else { - LOG_D("Path '" << realPath << "' is not a topic: " << entry.Kind); + YDB_LOG_DEBUG("Path is not a", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}, + {"topic", entry.Kind}); if (Settings.UserToken && !entry.SecurityObject->CheckAccess(NACLib::EAccessRights::DescribeSchema, *Settings.UserToken)) { - LOG_D("Path '" << realPath << "' UNAUTHORIZED"); + YDB_LOG_DEBUG("Path UNAUTHORIZED", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); Result[originalPath] = TTopicInfo{ .Status = EStatus::UNAUTHORIZED_WITH_DESCRIBE_ACCESS }; @@ -176,7 +201,9 @@ public: break; } default: { - LOG_D("Path '" << realPath << "' unknown error"); + YDB_LOG_DEBUG("Path unknown error", + {"logPrefix", LOG_PREFIX}, + {"realPath", realPath}); Result[originalPath] = TTopicInfo{ .Status = EStatus::UNKNOWN_ERROR, .RealPath = realPath diff --git a/ydb/core/persqueue/public/fetcher/fetch_request_actor.cpp b/ydb/core/persqueue/public/fetcher/fetch_request_actor.cpp index 044e02c8ae1..05951b7ff2d 100644 --- a/ydb/core/persqueue/public/fetcher/fetch_request_actor.cpp +++ b/ydb/core/persqueue/public/fetcher/fetch_request_actor.cpp @@ -11,12 +11,9 @@ #include <ydb/library/actors/core/actor_bootstrapped.h> #include <ydb/library/actors/core/log.h> -#define LOG_PREFIX "[" << NActors::TlsActivationContext->AsActorContext().SelfID << "] " -#define LOG_E(stream) LOG_ERROR_S(*NActors::TlsActivationContext, NKikimrServices::PQ_FETCH_REQUEST, LOG_PREFIX << stream) -#define LOG_W(stream) LOG_WARN_S(*NActors::TlsActivationContext, NKikimrServices::PQ_FETCH_REQUEST, LOG_PREFIX << stream) -#define LOG_I(stream) LOG_INFO_S(*NActors::TlsActivationContext, NKikimrServices::PQ_FETCH_REQUEST, LOG_PREFIX << stream) -#define LOG_D(stream) LOG_DEBUG_S(*NActors::TlsActivationContext, NKikimrServices::PQ_FETCH_REQUEST, LOG_PREFIX << stream) -#define LOG(level, stream) LOG_LOG_S(*NActors::TlsActivationContext, level, NKikimrServices::PQ_FETCH_REQUEST, LOG_PREFIX << stream) +#define YDB_LOG_THIS_FILE_COMPONENT NKikimrServices::PQ_FETCH_REQUEST + +#define LOG_PREFIX TStringBuilder() << "[" << NActors::TlsActivationContext->AsActorContext().SelfID << "] " namespace NKikimr::NPQ { @@ -152,7 +149,9 @@ public: } void Bootstrap(const TActorContext& ctx) { - LOG_I("Fetch request actor boostrapped. Request is valid: " << (!Response)); + YDB_LOG_INFO("Fetch request actor boostrapped. Request is", + {"logPrefix", LOG_PREFIX}, + {"valid", (!Response)}); // handle error from constructor if (Response) { @@ -167,7 +166,8 @@ public: } void DescribeTopics(const TActorContext&) { - LOG_D("DescribeTopics"); + YDB_LOG_DEBUG("DescribeTopics", + {"logPrefix", LOG_PREFIX}); std::unordered_set<TString> topics; for (const auto& part : Settings.Partitions) { @@ -182,7 +182,8 @@ public: } void Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev, const TActorContext& ctx) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", LOG_PREFIX}); for (auto& [topicPath, info] : ev->Get()->Topics) { switch (info.Status) { @@ -245,7 +246,9 @@ public: auto& fetchInfo = topicInfo.HasDataRequests[partitionId]; const auto partitionIndex = fetchInfo->Record.GetCookie(); tabletInfo.PartitionIndexes.push_back(partitionIndex); - LOG_D("Sending TEvPersQueue::TEvHasDataInfo " << fetchInfo->Record.ShortDebugString()); + YDB_LOG_DEBUG("Sending TEvPersQueue::TEvHasDataInfo", + {"logPrefix", LOG_PREFIX}, + {"fetchInfoRecord", fetchInfo->Record.ShortDebugString()}); NTabletPipe::SendData(ctx, tabletInfo.PipeClient, fetchInfo.Release()); PartitionStatus[partitionIndex] = EPartitionStatus::HasDataRequested; } @@ -260,8 +263,11 @@ public: void ProceedFetchRequest(const TActorContext& ctx) { if (FetchRequestCurrentReadTablet) { //already got active read request - LOG_D("Fetch request is pending. TabletId=" << FetchRequestCurrentReadTablet - << " partitionIndex=" << FetchRequestCurrentPartitionIndex << "/" << Settings.Partitions.size()); + YDB_LOG_DEBUG("Fetch request is pending. /", + {"logPrefix", LOG_PREFIX}, + {"tabletId", FetchRequestCurrentReadTablet}, + {"partitionIndex", FetchRequestCurrentPartitionIndex}, + {"settingsPartitionsSize", Settings.Partitions.size()}); return; } @@ -270,7 +276,10 @@ public: ("r", Settings.Partitions.size()); while (true) { - LOG_D("Processing " << FetchRequestCurrentPartitionIndex << "/" << Settings.Partitions.size()); + YDB_LOG_DEBUG("Processing /", + {"logPrefix", LOG_PREFIX}, + {"fetchRequestCurrentPartitionIndex", FetchRequestCurrentPartitionIndex}, + {"settingsPartitionsSize", Settings.Partitions.size()}); if (FetchRequestCurrentPartitionIndex == Settings.Partitions.size()) { CreateOkResponse(); return SendReplyAndDie(std::move(Response), ctx); @@ -278,13 +287,19 @@ public: auto& status = PartitionStatus[FetchRequestCurrentPartitionIndex]; if (status == EPartitionStatus::DataReceived) { - LOG_D("Skip partition " << FetchRequestCurrentPartitionIndex << " because status is DataReceived"); + YDB_LOG_DEBUG("Skip partition because status is DataReceived", + {"logPrefix", LOG_PREFIX}, + {"fetchRequestCurrentPartitionIndex", FetchRequestCurrentPartitionIndex}); ++FetchRequestCurrentPartitionIndex; continue; } if (FetchRequestBytesLeft == 0) { - LOG_D("Partition " << FetchRequestCurrentPartitionIndex << " status is " << (int)status << " bytesLeft=" << FetchRequestBytesLeft); + YDB_LOG_DEBUG("Partition status is", + {"logPrefix", LOG_PREFIX}, + {"fetchRequestCurrentPartitionIndex", FetchRequestCurrentPartitionIndex}, + {"status", (int)status}, + {"bytesLeft", FetchRequestBytesLeft}); if (status == EPartitionStatus::HasDataReceived) { ++FetchRequestCurrentPartitionIndex; continue; @@ -316,7 +331,9 @@ public: //Form read request auto request = CreateReadRequest(topic, req); - LOG_D("Sending request: " << request->Record.ShortDebugString()); + YDB_LOG_DEBUG("Sending", + {"logPrefix", LOG_PREFIX}, + {"request", request->Record.ShortDebugString()}); NTabletPipe::SendData(ctx, tabletInfo.PipeClient, request.release()); break; @@ -326,13 +343,18 @@ public: void Handle(TEvPersQueue::TEvHasDataInfoResponse::TPtr& ev, const TActorContext& ctx) { auto& record = ev->Get()->Record; auto partitionIndex = record.GetCookie(); - LOG_D("Handle TEvPersQueue::TEvHasDataInfoResponse " << record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPersQueue::TEvHasDataInfoResponse", + {"logPrefix", LOG_PREFIX}, + {"ev", record.ShortDebugString()}); if (partitionIndex >= PartitionStatus.size()) { Y_VERIFY_DEBUG(partitionIndex < PartitionStatus.size()); return; } auto& status = PartitionStatus[partitionIndex]; - LOG_D("Partition " << partitionIndex << " status is " << (int)status); + YDB_LOG_DEBUG("Partition status is", + {"logPrefix", LOG_PREFIX}, + {"partitionIndex", partitionIndex}, + {"status", (int)status}); if (status != EPartitionStatus::HasDataRequested) { // On timeout we resend send HasData return; @@ -354,12 +376,17 @@ public: void Handle(TEvPersQueue::TEvResponse::TPtr& ev, const TActorContext& ctx) { auto& record = ev->Get()->Record; - LOG_D("Handle TEvPersQueue::TEvResponse " << record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPersQueue::TEvResponse", + {"logPrefix", LOG_PREFIX}, + {"ev", record.ShortDebugString()}); AFL_ENSURE(record.HasPartitionResponse()); if (record.GetPartitionResponse().GetCookie() != FetchRequestCurrentPartitionIndex || FetchRequestCurrentReadTablet == 0) { - LOG_W("proxy fetch error: got response from tablet " << record.GetPartitionResponse().GetCookie() - << " while waiting from " << FetchRequestCurrentPartitionIndex << " and requested tablet is " << FetchRequestCurrentReadTablet); + YDB_LOG_WARN("Proxy fetch error: got response from tablet while waiting from and requested tablet is", + {"logPrefix", LOG_PREFIX}, + {"partitionResponseCookie", record.GetPartitionResponse().GetCookie()}, + {"fetchRequestCurrentPartitionIndex", FetchRequestCurrentPartitionIndex}, + {"fetchRequestCurrentReadTablet", FetchRequestCurrentReadTablet}); return; } @@ -401,7 +428,9 @@ public: Response->Response.SetTimestampType(timestampType); } - LOG_D("After processing result FetchRequestBytesLeft=" << FetchRequestBytesLeft); + YDB_LOG_DEBUG("After processing result", + {"logPrefix", LOG_PREFIX}, + {"fetchRequestBytesLeft", FetchRequestBytesLeft}); if (FetchRequestBytesLeft == 0) { FinishProcessing(ctx); } else if (IsQuotaRequired()) { @@ -489,7 +518,9 @@ public: fetchInfo->Record.SetDeadline(0); auto tabletId = topicInfo.PartitionToTablet[p.Partition]; - LOG_D("Sending TEvPersQueue::TEvHasDataInfo " << fetchInfo->Record.ShortDebugString()); + YDB_LOG_DEBUG("Sending TEvPersQueue::TEvHasDataInfo", + {"logPrefix", LOG_PREFIX}, + {"fetchInfoRecord", fetchInfo->Record.ShortDebugString()}); NTabletPipe::SendData(ctx, TabletInfo[tabletId].PipeClient, fetchInfo.release()); } } @@ -541,7 +572,10 @@ public: } void SendReplyAndDie(THolder<TEvPQ::TEvFetchResponse> event, const TActorContext& ctx) { - LOG_D("Reply to " << RequesterId << ": " << event->Response.ShortDebugString()); + YDB_LOG_DEBUG("Reply", + {"logPrefix", LOG_PREFIX}, + {"requesterId", RequesterId}, + {"response", event->Response.ShortDebugString()}); ctx.Send(RequesterId, event.Release()); Die(ctx); } diff --git a/ydb/core/persqueue/public/mlp/mlp_changer.h b/ydb/core/persqueue/public/mlp/mlp_changer.h index 75a3c949b87..9ecb14a68c5 100644 --- a/ydb/core/persqueue/public/mlp/mlp_changer.h +++ b/ydb/core/persqueue/public/mlp/mlp_changer.h @@ -48,7 +48,8 @@ public: private: void DoDescribe() { - LOG_D("Start describe"); + YDB_LOG_DEBUG_COMP(Service, "Start describe", + {"logPrefix", NPQ_LOG_PREFIX}); TBase::Become(&TThis::DescribeState); NDescriber::TDescribeSettings settings = { @@ -59,7 +60,8 @@ private: } void Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG_COMP(Service, "Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); ChildActorId = {}; @@ -93,7 +95,8 @@ private: } void DoChanges() { - LOG_D("Start DoChanges"); + YDB_LOG_DEBUG_COMP(Service, "Start DoChanges", + {"logPrefix", NPQ_LOG_PREFIX}); TBase::Become(&TThis::ChangesState); for (const TMessageId& messageId: Settings.Messages) { @@ -121,12 +124,16 @@ private: } void Handle(typename TResponse::TPtr& ev) { - LOG_D("Handle response " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG_COMP(Service, "Handle response", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto partitionId = ev->Cookie; auto it = PendingPartitions.find(partitionId); if (it == PendingPartitions.end()) { - LOG_D("Received response fron unexpected partition " << partitionId); + YDB_LOG_DEBUG_COMP(Service, "Received response fron unexpected partition", + {"logPrefix", NPQ_LOG_PREFIX}, + {"partitionId", partitionId}); return; } @@ -144,13 +151,17 @@ private: } void Handle(TEvPQ::TEvMLPErrorResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPErrorResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG_COMP(Service, "Handle TEvPQ::TEvMLPErrorResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto partitionId = ev->Cookie; auto it = PendingPartitions.find(partitionId); if (it == PendingPartitions.end()) { - LOG_D("Received response from unexpected partition " << partitionId); + YDB_LOG_DEBUG_COMP(Service, "Received response from unexpected partition", + {"logPrefix", NPQ_LOG_PREFIX}, + {"partitionId", partitionId}); return; } @@ -164,11 +175,15 @@ private: } void Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem " << ev->Get()->TabletId); + YDB_LOG_DEBUG_COMP(Service, "Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}, + {"TabletId", ev->Get()->TabletId}); auto it = Pipes.find(ev->Get()->TabletId); if (it == Pipes.end()) { - LOG_D("Received pipe error for unexpected tablet " << ev->Get()->TabletId); + YDB_LOG_DEBUG_COMP(Service, "Received pipe error for unexpected tablet", + {"logPrefix", NPQ_LOG_PREFIX}, + {"TabletId", ev->Get()->TabletId}); return; } @@ -241,7 +256,9 @@ private: } void ReplyErrorAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_I("Reply error " << Ydb::StatusIds::StatusCode_Name(errorCode)); + YDB_LOG_INFO_COMP(Service, "Reply error", + {"logPrefix", NPQ_LOG_PREFIX}, + {"statusCodeName", Ydb::StatusIds::StatusCode_Name(errorCode)}); TBase::Send(ParentId, new TEvChangeResponse(errorCode, std::move(errorMessage))); PassAway(); } diff --git a/ydb/core/persqueue/public/mlp/mlp_describer.cpp b/ydb/core/persqueue/public/mlp/mlp_describer.cpp index b2203a7d6cb..8d54944c301 100644 --- a/ydb/core/persqueue/public/mlp/mlp_describer.cpp +++ b/ydb/core/persqueue/public/mlp/mlp_describer.cpp @@ -3,6 +3,8 @@ #include <ydb/core/persqueue/public/constants.h> #include <ydb/core/persqueue/public/utils.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NMLP { TDescriberActor::TDescriberActor(const TActorId& parentId, const TDescribeSettings& settings) @@ -17,7 +19,8 @@ void TDescriberActor::Bootstrap() { } void TDescriberActor::DoDescribe() { - LOG_D("Start describe"); + YDB_LOG_DEBUG("Start describe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TDescriberActor::DescribeState); NDescriber::TDescribeSettings settings = { @@ -28,7 +31,8 @@ void TDescriberActor::DoDescribe() { } void TDescriberActor::Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); ChildActorId = {}; @@ -61,13 +65,16 @@ STFUNC(TDescriberActor::DescribeState) { } void TDescriberActor::DoRuntimeAttributes() { - LOG_D("Start DoRuntimeAttributes"); + YDB_LOG_DEBUG("Start DoRuntimeAttributes", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TDescriberActor::RuntimeAttributesState); SendToTablet(TopicInfo.Info->Description.GetBalancerTabletID(), new TEvPQ::TEvMLPGetRuntimeAttributesRequest(Settings.TopicName, Settings.Consumer)); } void TDescriberActor::Handle(TEvPQ::TEvMLPGetRuntimeAttributesResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPGetRuntimeAttributesResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPGetRuntimeAttributesResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto* result = ev->Get(); auto response = std::make_unique<TEvDescribeResponse>(); @@ -84,7 +91,8 @@ void TDescriberActor::Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { if (ev->Cookie != Cookie) { return; } - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); if (Backoff.HasMore()) { Backoff.Next(); return DoRuntimeAttributes(); @@ -103,7 +111,9 @@ STFUNC(TDescriberActor::RuntimeAttributesState) { void TDescriberActor::Handle(TEvPQ::TEvMLPErrorResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPErrorResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPErrorResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); ReplyErrorAndDie(ev->Get()->GetStatus(), std::move(ev->Get()->GetErrorMessage())); } @@ -113,7 +123,9 @@ void TDescriberActor::SendToTablet(ui64 tabletId, IEventBase *ev) { } void TDescriberActor::ReplyErrorAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_I("Reply error " << Ydb::StatusIds::StatusCode_Name(errorCode)); + YDB_LOG_INFO("Reply error", + {"logPrefix", NPQ_LOG_PREFIX}, + {"statusCodeName", Ydb::StatusIds::StatusCode_Name(errorCode)}); Send(ParentId, new TEvDescribeResponse(errorCode, std::move(errorMessage))); PassAway(); } diff --git a/ydb/core/persqueue/public/mlp/mlp_purger.cpp b/ydb/core/persqueue/public/mlp/mlp_purger.cpp index 17ffcb8e257..4f58897866e 100644 --- a/ydb/core/persqueue/public/mlp/mlp_purger.cpp +++ b/ydb/core/persqueue/public/mlp/mlp_purger.cpp @@ -1,5 +1,7 @@ #include "mlp_purger.h" +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NMLP { TPurgerActor::TPurgerActor(const TActorId& parentId, const TPurgerSettings& settings) @@ -14,7 +16,8 @@ void TPurgerActor::Bootstrap() { } void TPurgerActor::DoDescribe() { - LOG_D("Start describe"); + YDB_LOG_DEBUG("Start describe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TPurgerActor::DescribeState); NDescriber::TDescribeSettings settings = { @@ -25,7 +28,8 @@ void TPurgerActor::DoDescribe() { } void TPurgerActor::Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); ChildActorId = {}; @@ -53,7 +57,8 @@ STFUNC(TPurgerActor::DescribeState) { } void TPurgerActor::DoPurge() { - LOG_D("Start purge"); + YDB_LOG_DEBUG("Start purge", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TPurgerActor::PurgeState); for (auto& partition : TopicInfo.Info->Description.GetPartitions()) { @@ -69,7 +74,9 @@ void TPurgerActor::DoPurge() { void TPurgerActor::Handle(TEvPQ::TEvMLPPurgeResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPPurgeResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPPurgeResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto partitionId = ev->Get()->GetPartitionId(); auto& partitionStatus = Partitions[partitionId]; @@ -100,7 +107,9 @@ void TPurgerActor::RetryIfPossible(ui32 partitionId, TPartitionStatus& partition void TPurgerActor::Handle(TEvPQ::TEvMLPErrorResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPErrorResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPErrorResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto partitionId = ev->Get()->GetPartitionId(); auto& partitionStatus = Partitions[partitionId]; @@ -118,7 +127,8 @@ void TPurgerActor::Handle(TEvPQ::TEvMLPErrorResponse::TPtr& ev) void TPurgerActor::Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); auto tabletId = ev->Get()->TabletId; ++TabletCookies[tabletId]; @@ -133,7 +143,8 @@ void TPurgerActor::Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) } void TPurgerActor::Handle(TEvents::TEvWakeup::TPtr& ev) { - LOG_D("Handle TEvents::TEvWakeup"); + YDB_LOG_DEBUG("Handle TEvents::TEvWakeup", + {"logPrefix", NPQ_LOG_PREFIX}); auto partitionId = ev->Get()->Tag; auto& partitionStatus = Partitions[partitionId]; @@ -168,7 +179,10 @@ void TPurgerActor::RequestPartitionIfNeeded(ui32 partitionId,TPartitionStatus& s } void TPurgerActor::ReplyIfPossible() { - LOG_D("ReplyIfPossible: PendingPartitions " << PendingPartitions << " PendingRetries " << PendingRetries); + YDB_LOG_DEBUG("ReplyIfPossible: PendingPartitions PendingRetries", + {"logPrefix", NPQ_LOG_PREFIX}, + {"pendingPartitions", PendingPartitions}, + {"pendingRetries", PendingRetries}); if (PendingPartitions > 0 || PendingRetries > 0) { return; } @@ -190,7 +204,9 @@ void TPurgerActor::SendToTablet(ui64 tabletId, IEventBase *ev, ui64 cookie) { } void TPurgerActor::ReplyErrorAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_I("Reply error " << Ydb::StatusIds::StatusCode_Name(errorCode)); + YDB_LOG_INFO("Reply error", + {"logPrefix", NPQ_LOG_PREFIX}, + {"statusCodeName", Ydb::StatusIds::StatusCode_Name(errorCode)}); Send(ParentId, new TEvPurgeResponse(errorCode, std::move(errorMessage))); PassAway(); } diff --git a/ydb/core/persqueue/public/mlp/mlp_reader.cpp b/ydb/core/persqueue/public/mlp/mlp_reader.cpp index 0cf8a4a7221..0ebe3caf0f8 100644 --- a/ydb/core/persqueue/public/mlp/mlp_reader.cpp +++ b/ydb/core/persqueue/public/mlp/mlp_reader.cpp @@ -7,6 +7,8 @@ #include <ydb/public/api/protos/ydb_topic.pb.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/topic/codecs.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NMLP { TReaderActor::TReaderActor(const TActorId& parentId, const TReaderSettings& settings) @@ -21,7 +23,8 @@ void TReaderActor::Bootstrap() { } void TReaderActor::DoDescribe() { - LOG_D("Start describe"); + YDB_LOG_DEBUG("Start describe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TReaderActor::DescribeState); NDescriber::TDescribeSettings settings = { @@ -32,7 +35,8 @@ void TReaderActor::DoDescribe() { } void TReaderActor::Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); ChildActorId = {}; @@ -65,13 +69,16 @@ STFUNC(TReaderActor::DescribeState) { } void TReaderActor::DoSelectPartition() { - LOG_D("Start select partition"); + YDB_LOG_DEBUG("Start select partition", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TReaderActor::SelectPartitionState); SendToTablet(Info->Description.GetBalancerTabletID(), new TEvPQ::TEvMLPGetPartitionRequest(Settings.TopicName, Settings.Consumer, Settings.ReceiveAttemptId)); } void TReaderActor::Handle(TEvPQ::TEvMLPGetPartitionResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPGetPartitionResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPGetPartitionResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); auto* result = ev->Get(); switch (result->GetStatus()) { case Ydb::StatusIds::SUCCESS: { @@ -88,7 +95,8 @@ void TReaderActor::HandleOnSelectPartition(TEvPipeCache::TEvDeliveryProblem::TPt if (ev->Cookie != Cookie) { return; } - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); if (Backoff.HasMore()) { Backoff.Next(); return DoSelectPartition(); @@ -106,7 +114,8 @@ STFUNC(TReaderActor::SelectPartitionState) { } void TReaderActor::DoRead() { - LOG_D("Start read"); + YDB_LOG_DEBUG("Start read", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TReaderActor::ReadState); auto* request = new TEvPQ::TEvMLPReadRequest( @@ -123,7 +132,8 @@ void TReaderActor::DoRead() { } void TReaderActor::Handle(TEvPQ::TEvMLPReadResponse::TPtr& ev) { - LOG_D("Handle TEvPQ::TEvMLPReadResponse"); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPReadResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto response = std::make_unique<TEvReadResponse>(); response->BalancerTabletId = Info->Description.GetBalancerTabletID(); @@ -131,7 +141,9 @@ void TReaderActor::Handle(TEvPQ::TEvMLPReadResponse::TPtr& ev) { NKikimrPQClient::TDataChunk proto; bool res = proto.ParseFromString(message.GetData()); if (!res) { - LOG_W("Error parsing data. Offset " << message.GetId().GetOffset()); + YDB_LOG_WARN("Error parsing data. Offset", + {"logPrefix", NPQ_LOG_PREFIX}, + {"messageIdOffset", message.GetId().GetOffset()}); // Skip message continue; } @@ -173,7 +185,9 @@ void TReaderActor::Handle(TEvPQ::TEvMLPReadResponse::TPtr& ev) { void TReaderActor::Handle(TEvPQ::TEvMLPErrorResponse::TPtr& ev) { // TODO MLP Retry - LOG_D("Handle TEvPQ::TEvMLPErrorResponse " << ev->Get()->Record.ShortDebugString()); + YDB_LOG_DEBUG("Handle TEvPQ::TEvMLPErrorResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ev", ev->Get()->Record.ShortDebugString()}); ReplyErrorAndDie(ev->Get()->GetStatus(), std::move(ev->Get()->GetErrorMessage())); } @@ -181,7 +195,8 @@ void TReaderActor::HandleOnRead(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { if (ev->Cookie != Cookie) { return; } - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); if (Backoff.HasMore()) { Backoff.Next(); return DoRead(); @@ -205,7 +220,9 @@ void TReaderActor::SendToTablet(ui64 tabletId, IEventBase *ev) { } void TReaderActor::ReplyErrorAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_I("Reply error " << Ydb::StatusIds::StatusCode_Name(errorCode)); + YDB_LOG_INFO("Reply error", + {"logPrefix", NPQ_LOG_PREFIX}, + {"statusCodeName", Ydb::StatusIds::StatusCode_Name(errorCode)}); Send(ParentId, new TEvReadResponse(errorCode, std::move(errorMessage))); PassAway(); } diff --git a/ydb/core/persqueue/public/mlp/mlp_writer.cpp b/ydb/core/persqueue/public/mlp/mlp_writer.cpp index beff325c438..5c05b6a9ca1 100644 --- a/ydb/core/persqueue/public/mlp/mlp_writer.cpp +++ b/ydb/core/persqueue/public/mlp/mlp_writer.cpp @@ -5,6 +5,8 @@ #include <ydb/public/api/protos/ydb_topic.pb.h> #include <ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/topic/codecs.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NMLP { TWriterActor::TWriterActor(const TActorId& parentId, const TWriterSettings& settings) @@ -27,7 +29,8 @@ void TWriterActor::PassAway() { } void TWriterActor::DoDescribe() { - LOG_D("Start describe"); + YDB_LOG_DEBUG("Start describe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TWriterActor::DescribeState); NDescriber::TDescribeSettings settings = { @@ -38,7 +41,8 @@ void TWriterActor::DoDescribe() { } void TWriterActor::Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); ChildActorId = {}; @@ -118,7 +122,8 @@ size_t SerializeTo(TWriterSettings::TMessage& item, ::NKikimrClient::TPersQueueP } void TWriterActor::DoWrite() { - LOG_D("Start write"); + YDB_LOG_DEBUG("Start write", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TWriterActor::WriteState); struct TInfo { @@ -177,7 +182,8 @@ void TWriterActor::DoWrite() { } void TWriterActor::Handle(TEvPersQueue::TEvResponse::TPtr& ev) { - LOG_D("Handle TEvPersQueue::TEvResponse"); + YDB_LOG_DEBUG("Handle TEvPersQueue::TEvResponse", + {"logPrefix", NPQ_LOG_PREFIX}); bool alreadyReceived = false; auto& record = ev->Get()->Record; @@ -210,7 +216,8 @@ void TWriterActor::Handle(TEvPersQueue::TEvResponse::TPtr& ev) { } void TWriterActor::Handle(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); const auto tabletId = ev->Get()->TabletId; @@ -242,8 +249,12 @@ void TWriterActor::SendToTablet(ui64 tabletId, IEventBase *ev) { } bool TWriterActor::OnUnhandledException(const std::exception& exc) { - LOG_C("unhandled exception " << TypeName(exc) << ": " << exc.what() << Endl - << TBackTrace::FromCurrentException().PrintToString()); + YDB_LOG_CRIT("Unhandled exception", + {"logPrefix", NPQ_LOG_PREFIX}, + {"exceptionType", TypeName(exc)}, + {"exceptionMessage", exc.what()}, + {"endl", Endl}, + {"backTrace", TBackTrace::FromCurrentException().PrintToString()}); PendingRequests = 0; ReplyIfPossible(); @@ -253,11 +264,15 @@ bool TWriterActor::OnUnhandledException(const std::exception& exc) { bool TWriterActor::IsSuccess(const NKikimrClient::TResponse& record) { if (record.HasErrorCode() && record.GetErrorCode() != NPersQueue::NErrorCode::OK) { - LOG_W("Write error: " << record.ShortDebugString()); + YDB_LOG_WARN("Write", + {"logPrefix", NPQ_LOG_PREFIX}, + {"error", record.ShortDebugString()}); return false; } if (!record.HasPartitionResponse()) { - LOG_W("Missing partition response: " << record.ShortDebugString()); + YDB_LOG_WARN("Missing partition", + {"logPrefix", NPQ_LOG_PREFIX}, + {"response", record.ShortDebugString()}); return false; } diff --git a/ydb/core/persqueue/public/schema/alter_topic_internal.cpp b/ydb/core/persqueue/public/schema/alter_topic_internal.cpp index e41f79b1752..8b235fafc0d 100644 --- a/ydb/core/persqueue/public/schema/alter_topic_internal.cpp +++ b/ydb/core/persqueue/public/schema/alter_topic_internal.cpp @@ -4,6 +4,8 @@ #include <ydb/services/persqueue_v1/actors/events.h> #include <ydb/services/persqueue_v1/actors/schema/common/grpc_proxy_actor.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -30,7 +32,9 @@ public: } void OnException(const std::exception& exc) override { - LOG_E("OnException: " << exc.what()); + YDB_LOG_ERROR("Catch exception", + {"logPrefix", NPQ_LOG_PREFIX}, + {"onException", exc.what()}); TEvSchemaResponse response(Path, Ydb::StatusIds::INTERNAL_ERROR, exc.what()); @@ -43,7 +47,10 @@ public: private: void Handle(NPQ::NSchema::TEvSchemaResponse::TPtr& ev) { - LOG_D("Handle TEvSchemaResponse. Status: " << ev->Get()->Status << ", ErrorMessage: " << ev->Get()->ErrorMessage); + YDB_LOG_DEBUG("Handle TEvSchemaResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"status", ev->Get()->Status}, + {"errorMessage", ev->Get()->ErrorMessage}); Promise.SetValue({ .Path = Path, diff --git a/ydb/core/persqueue/public/schema/alter_topic_operation.cpp b/ydb/core/persqueue/public/schema/alter_topic_operation.cpp index bc1f6080ecd..18ddc3dae9d 100644 --- a/ydb/core/persqueue/public/schema/alter_topic_operation.cpp +++ b/ydb/core/persqueue/public/schema/alter_topic_operation.cpp @@ -6,6 +6,8 @@ #include <ydb/core/protos/schemeshard/operations.pb.h> #include <ydb/core/ydb_convert/tx_proxy_status.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -37,7 +39,8 @@ public: private: void DoDescribe() { - LOG_D("DoDescribe"); + YDB_LOG_DEBUG("DoDescribe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TAlterTopicOperationActor::DescribeState); RegisterWithSameMailbox(NDescriber::CreateDescriberActor( @@ -52,7 +55,8 @@ private: } void Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& topics = ev->Get()->Topics; AFL_ENSURE(topics.size() == 1)("s", topics.size()); @@ -90,14 +94,16 @@ private: private: void DoGetClustersList() { - LOG_D("DoGetClustersList"); + YDB_LOG_DEBUG("DoGetClustersList", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TAlterTopicOperationActor::GetClustersListState); Send(NPQ::NClusterTracker::MakeClusterTrackerID(), new NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersList()); } void Handle(NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersListResponse::TPtr& ev) { - LOG_D("Handle NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersListResponse: " - << (ev->Get()->Success ? ev->Get()->ClustersList->DebugString() : "error")); + YDB_LOG_DEBUG("Handle", + {"logPrefix", NPQ_LOG_PREFIX}, + {"getClustersListResponse", (ev->Get()->Success ? ev->Get()->ClustersList->DebugString() : "error")}); auto& response = *ev->Get(); if (response.Success) { @@ -116,7 +122,8 @@ private: private: void DoAlter() { - LOG_D("DoAlter"); + YDB_LOG_DEBUG("DoAlter", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TAlterTopicOperationActor::AlterState); @@ -178,7 +185,8 @@ private: } void Handle(TEvSchemaOperationResponse::TPtr& ev) { - LOG_D("Handle TEvSchemaOperationResponse"); + YDB_LOG_DEBUG("Handle TEvSchemaOperationResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& response = *ev->Get(); return ReplyAndDie(response.Status, std::move(response.ErrorMessage)); } @@ -192,7 +200,10 @@ private: private: void ReplyAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_D("ReplyAndDie " << errorCode << " '" << errorMessage << "'"); + YDB_LOG_DEBUG("ReplyAndDie", + {"logPrefix", NPQ_LOG_PREFIX}, + {"errorCode", errorCode}, + {"errorMessage", errorMessage}); if (errorCode == Ydb::StatusIds::SUCCESS && !Settings.PrepareOnly) { ModifyScheme = {}; } diff --git a/ydb/core/persqueue/public/schema/create_topic_internal.cpp b/ydb/core/persqueue/public/schema/create_topic_internal.cpp index dc7ae9d63be..2d7df05a30a 100644 --- a/ydb/core/persqueue/public/schema/create_topic_internal.cpp +++ b/ydb/core/persqueue/public/schema/create_topic_internal.cpp @@ -4,6 +4,8 @@ #include <ydb/services/persqueue_v1/actors/events.h> #include <ydb/services/persqueue_v1/actors/schema/common/grpc_proxy_actor.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -30,7 +32,9 @@ public: } void OnException(const std::exception& exc) override { - LOG_E("OnException: " << exc.what()); + YDB_LOG_ERROR("Catch exception", + {"logPrefix", NPQ_LOG_PREFIX}, + {"onException", exc.what()}); TEvSchemaResponse response(Path, Ydb::StatusIds::INTERNAL_ERROR, exc.what()); @@ -43,7 +47,10 @@ public: private: void Handle(NPQ::NSchema::TEvSchemaResponse::TPtr& ev) { - LOG_D("Handle TEvSchemaResponse. Status: " << ev->Get()->Status << ", ErrorMessage: " << ev->Get()->ErrorMessage); + YDB_LOG_DEBUG("Handle TEvSchemaResponse", + {"logPrefix", NPQ_LOG_PREFIX}, + {"status", ev->Get()->Status}, + {"errorMessage", ev->Get()->ErrorMessage}); Promise.SetValue({ .Path = Path, diff --git a/ydb/core/persqueue/public/schema/create_topic_operation.cpp b/ydb/core/persqueue/public/schema/create_topic_operation.cpp index d24ad5194ed..a257c50b78c 100644 --- a/ydb/core/persqueue/public/schema/create_topic_operation.cpp +++ b/ydb/core/persqueue/public/schema/create_topic_operation.cpp @@ -8,6 +8,8 @@ #include <ydb/core/protos/schemeshard/operations.pb.h> #include <ydb/core/ydb_convert/tx_proxy_status.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -42,13 +44,15 @@ public: private: void DoGetClustersList() { - LOG_D("DoGetClustersList"); + YDB_LOG_DEBUG("DoGetClustersList", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TCreateTopicOperationActor::GetClustersListState); Send(NPQ::NClusterTracker::MakeClusterTrackerID(), new NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersList()); } void Handle(NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersListResponse::TPtr& ev) { - LOG_D("Handle NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersListResponse"); + YDB_LOG_DEBUG("Handle NPQ::NClusterTracker::TEvClusterTracker::TEvGetClustersListResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& response = *ev->Get(); if (response.Success) { @@ -67,7 +71,9 @@ private: private: void DoCreate() { - LOG_D("DoCreate IfNotExists: " << Settings.IfNotExists); + YDB_LOG_DEBUG("DoCreate", + {"logPrefix", NPQ_LOG_PREFIX}, + {"ifNotExists", Settings.IfNotExists}); Become(&TCreateTopicOperationActor::CreateState); auto database = CanonizePath(Settings.Database); @@ -117,7 +123,8 @@ private: } void Handle(TEvSchemaOperationResponse::TPtr& ev) { - LOG_D("Handle TEvSchemaOperationResponse"); + YDB_LOG_DEBUG("Handle TEvSchemaOperationResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& response = *ev->Get(); return ReplyAndDie(response.Status, std::move(response.ErrorMessage)); } @@ -131,7 +138,10 @@ private: private: void ReplyAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_D("ReplyAndDie " << errorCode << " '" << errorMessage << "'"); + YDB_LOG_DEBUG("ReplyAndDie", + {"logPrefix", NPQ_LOG_PREFIX}, + {"errorCode", errorCode}, + {"errorMessage", errorMessage}); if ((errorCode == Ydb::StatusIds::SUCCESS || errorCode == Ydb::StatusIds::ALREADY_EXISTS) && !Settings.PrepareOnly) { ModifyScheme = {}; } diff --git a/ydb/core/persqueue/public/schema/drop_topic_operation.cpp b/ydb/core/persqueue/public/schema/drop_topic_operation.cpp index ea7702a8c87..61e723667dc 100644 --- a/ydb/core/persqueue/public/schema/drop_topic_operation.cpp +++ b/ydb/core/persqueue/public/schema/drop_topic_operation.cpp @@ -8,6 +8,8 @@ #include <ydb/core/tx/schemeshard/schemeshard.h> #include <ydb/core/tx/tx_proxy/proxy.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -39,7 +41,8 @@ public: private: void DoDescribe() { - LOG_D("DoDescribe"); + YDB_LOG_DEBUG("DoDescribe", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TDropTopicOperationActor::DescribeState); RegisterWithSameMailbox(NDescriber::CreateDescriberActor( @@ -54,7 +57,8 @@ private: } void Handle(NDescriber::TEvDescribeTopicsResponse::TPtr& ev) { - LOG_D("Handle NDescriber::TEvDescribeTopicsResponse"); + YDB_LOG_DEBUG("Handle NDescriber::TEvDescribeTopicsResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& topics = ev->Get()->Topics; AFL_ENSURE(topics.size() == 1)("s", topics.size()); @@ -91,7 +95,8 @@ private: private: void DoDrop() { - LOG_D("DoDrop"); + YDB_LOG_DEBUG("DoDrop", + {"logPrefix", NPQ_LOG_PREFIX}); Become(&TDropTopicOperationActor::DropState); @@ -126,7 +131,8 @@ private: } void Handle(TEvSchemaOperationResponse::TPtr& ev) { - LOG_D("Handle TEvSchemaOperationResponse"); + YDB_LOG_DEBUG("Handle TEvSchemaOperationResponse", + {"logPrefix", NPQ_LOG_PREFIX}); auto& response = *ev->Get(); return ReplyAndDie(response.Status, std::move(response.ErrorMessage)); } diff --git a/ydb/core/persqueue/public/schema/schema_operation.cpp b/ydb/core/persqueue/public/schema/schema_operation.cpp index 53c06a376f7..0f98bb13e42 100644 --- a/ydb/core/persqueue/public/schema/schema_operation.cpp +++ b/ydb/core/persqueue/public/schema/schema_operation.cpp @@ -10,6 +10,8 @@ #include <ydb/core/ydb_convert/tx_proxy_status.h> #include <ydb/library/services/services.pb.h> +#define YDB_LOG_THIS_FILE_COMPONENT Service + namespace NKikimr::NPQ::NSchema { namespace { @@ -55,7 +57,9 @@ public: private: void DoPropose() { - LOG_D("DoPropose retry: " << ProposeBackoff.GetIteration()); + YDB_LOG_DEBUG("DoPropose", + {"logPrefix", NPQ_LOG_PREFIX}, + {"retry", ProposeBackoff.GetIteration()}); Become(&TSchemaOperationActor::ProposeState); auto request = std::make_unique<TEvTxUserProxy::TEvProposeTransaction>(); @@ -76,7 +80,8 @@ private: } void Handle(TEvTxUserProxy::TEvProposeTransactionStatus::TPtr& ev) { - LOG_D("Handle TEvTxUserProxy::TEvProposeTransactionStatus"); + YDB_LOG_DEBUG("Handle TEvTxUserProxy::TEvProposeTransactionStatus", + {"logPrefix", NPQ_LOG_PREFIX}); const auto status = ev->Get()->Status(); const auto& record = ev->Get()->Record; @@ -110,7 +115,8 @@ private: } void HandleOnPropose(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - LOG_D("HandleOnPropose TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("HandleOnPropose TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); if (TPipeCacheClient::OnUndelivered(ev)) { return ReplyErrorAndDie(Ydb::StatusIds::UNAVAILABLE, TStringBuilder() << "SchemeShard " << ev->Get()->TabletId << " is unavailable"); @@ -129,7 +135,10 @@ private: private: void DoWaitCompletion() { - LOG_D("DoWaitTxCompletion SchemeShardTabletId: " << SchemeShardTabletId << " TxId: " << TxId); + YDB_LOG_DEBUG("DoWaitTxCompletion", + {"logPrefix", NPQ_LOG_PREFIX}, + {"schemeShardTabletId", SchemeShardTabletId}, + {"txId", TxId}); Become(&TSchemaOperationActor::WaitCompletionState); auto request = std::make_unique<NSchemeShard::TEvSchemeShard::TEvNotifyTxCompletion>(TxId); @@ -137,12 +146,14 @@ private: } void Handle(NSchemeShard::TEvSchemeShard::TEvNotifyTxCompletionResult::TPtr&) { - LOG_D("Handle TEvSchemeShard::TEvNotifyTxCompletionResult"); + YDB_LOG_DEBUG("Handle TEvSchemeShard::TEvNotifyTxCompletionResult", + {"logPrefix", NPQ_LOG_PREFIX}); ReplyOkAndDie(); } void HandleOnWaitCompletion(TEvPipeCache::TEvDeliveryProblem::TPtr& ev) { - LOG_D("Handle TEvPipeCache::TEvDeliveryProblem"); + YDB_LOG_DEBUG("Handle TEvPipeCache::TEvDeliveryProblem", + {"logPrefix", NPQ_LOG_PREFIX}); OnUndelivered(ev); if (++WaitTxCompletionRetries > MaxWaitTxCompletionRetries) { return ReplyErrorAndDie(Ydb::StatusIds::UNAVAILABLE, @@ -166,13 +177,17 @@ private: } void ReplyErrorAndDie(Ydb::StatusIds::StatusCode errorCode, TString&& errorMessage) { - LOG_D("ReplyErrorAndDie: " << errorCode << " " << errorMessage); + YDB_LOG_DEBUG("Dump NPQLOGPREFIX, replyErrorAndDie, errorMessage", + {"logPrefix", NPQ_LOG_PREFIX}, + {"replyErrorAndDie", errorCode}, + {"errorMessage", errorMessage}); Send(ParentId, new TEvSchemaOperationResponse(errorCode, std::move(errorMessage)), 0, Cookie); PassAway(); } void ReplyOkAndDie() { - LOG_D("ReplyOkAndDie"); + YDB_LOG_DEBUG("ReplyOkAndDie", + {"logPrefix", NPQ_LOG_PREFIX}); Send(ParentId, new TEvSchemaOperationResponse(), 0, Cookie); PassAway(); } |
