summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorkseleznyov <[email protected]>2026-07-15 17:59:21 +0300
committerGitHub <[email protected]>2026-07-15 17:59:21 +0300
commit4cd9d470a6cc9034a46b8f19300220bc045a67f7 (patch)
treea33727c3aaa81cb4a04d9e31db7a21909a41b395
parentf3db5fee3486b7d66987cb0f05c03c1be7a91489 (diff)
[YDB_LOG] Migrate ydb/core/persqueue/public (#45809)
-rw-r--r--ydb/core/persqueue/public/cloud_events/actor.cpp18
-rw-r--r--ydb/core/persqueue/public/cluster_tracker/cluster_tracker.cpp36
-rw-r--r--ydb/core/persqueue/public/describer/describer.cpp59
-rw-r--r--ydb/core/persqueue/public/fetcher/fetch_request_actor.cpp82
-rw-r--r--ydb/core/persqueue/public/mlp/mlp_changer.h37
-rw-r--r--ydb/core/persqueue/public/mlp/mlp_describer.cpp26
-rw-r--r--ydb/core/persqueue/public/mlp/mlp_purger.cpp34
-rw-r--r--ydb/core/persqueue/public/mlp/mlp_reader.cpp39
-rw-r--r--ydb/core/persqueue/public/mlp/mlp_writer.cpp33
-rw-r--r--ydb/core/persqueue/public/schema/alter_topic_internal.cpp11
-rw-r--r--ydb/core/persqueue/public/schema/alter_topic_operation.cpp27
-rw-r--r--ydb/core/persqueue/public/schema/create_topic_internal.cpp11
-rw-r--r--ydb/core/persqueue/public/schema/create_topic_operation.cpp20
-rw-r--r--ydb/core/persqueue/public/schema/drop_topic_operation.cpp14
-rw-r--r--ydb/core/persqueue/public/schema/schema_operation.cpp31
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();
}