summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorkseleznyov <[email protected]>2026-07-17 22:15:11 +0300
committerGitHub <[email protected]>2026-07-17 22:15:11 +0300
commitd9d7a527424e3c8a8ee2dad0d3c983bfae0343ca (patch)
tree98bfc9dd667dbd230b9c9e14b941e6a9e3049651
parent6f56721d2b8340e6d62967c01a0de557e46e7390 (diff)
[YDB_LOG] Migrate ydb/core/tablet/resource_broker (#46900)
Co-authored-by: Ilia Shakhov <[email protected]>
-rw-r--r--ydb/core/tablet/resource_broker.cpp181
1 files changed, 96 insertions, 85 deletions
diff --git a/ydb/core/tablet/resource_broker.cpp b/ydb/core/tablet/resource_broker.cpp
index 8b146e56ea8..4679552c027 100644
--- a/ydb/core/tablet/resource_broker.cpp
+++ b/ydb/core/tablet/resource_broker.cpp
@@ -8,6 +8,8 @@
#include <util/generic/hash.h>
#include <util/string/join.h>
+#define YDB_LOG_THIS_FILE_COMPONENT NKikimrServices::RESOURCE_BROKER
+
namespace NKikimr {
namespace NResourceBroker {
@@ -387,10 +389,11 @@ bool TScheduler::SubmitTask(const TEvResourceBroker::TTask &task,
auto &config = TaskConfig(task.Type);
TTaskPtr newTask = new TTask(task, client, Now, config.Counters);
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Submitted new %s task %s priority=%" PRIu64 " resources={%s}",
- config.Name.data(), newTask->GetIdString().data(), task.Priority,
- JoinSeq(", ", task.RequiredResources).data());
+ YDB_LOG_DEBUG_CTX(as, "Submitted new task",
+ {"configName", config.Name.data()},
+ {"taskId", newTask->GetIdString().data()},
+ {"priority", task.Priority},
+ {"requiredResources", JoinSeq(", ", task.RequiredResources).data()});
auto id = std::make_pair(client, task.TaskId);
if (!task.TaskId) {
@@ -399,13 +402,12 @@ bool TScheduler::SubmitTask(const TEvResourceBroker::TTask &task,
} while (Tasks.contains(id));
newTask->TaskId = id.second;
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Use ID %" PRIu64 " for submitted task",
- id.second);
+ YDB_LOG_DEBUG_CTX(as, "Assigned generated task id",
+ {"generatedTaskId", id.second});
} else if (Tasks.contains(id)) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "SubmitTask failed for task %" PRIu64 " to %s: task with the same ID has been already submitted",
- task.TaskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "SubmitTask failed: task with the same id was already submitted",
+ {"submittedTaskId", task.TaskId},
+ {"client", ToString(client)});
return false;
}
@@ -428,17 +430,19 @@ bool TScheduler::UpdateTask(ui64 taskId,
{
auto it = Tasks.find(std::make_pair(client, taskId));
if (it == Tasks.end()) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "UpdateTask failed for task %" PRIu64 " to %s: cannot update unknown task",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "UpdateTask failed: cannot update unknown task",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return false;
}
auto task = it->second;
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Update task %s (priority=%" PRIu64 " type=%s resources={%s} resubmit=%" PRIu32 ")",
- task->GetIdString().data(), priority, type.data(), JoinSeq(", ", requiredResources).data(),
- (ui32)resubmit);
+ YDB_LOG_DEBUG_CTX(as, "Updated task resources",
+ {"taskId", task->GetIdString().data()},
+ {"priority", priority},
+ {"type", type.data()},
+ {"requiredResources", JoinSeq(", ", requiredResources).data()},
+ {"resubmit", (ui32)resubmit});
auto queue = task->Queue;
queue->EraseTask(task, false, Now);
@@ -461,15 +465,15 @@ bool TScheduler::UpdateTaskCookie(ui64 taskId,
{
auto it = Tasks.find(std::make_pair(client, taskId));
if (it == Tasks.end()) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "UpdateTaskCookie failed for task %" PRIu64 " to %s: cannot update unknown task's cookie",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "UpdateTaskCookie failed: cannot update unknown task's cookie",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return false;
}
auto task = it->second;
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Update cookie for task %s", task->GetIdString().data());
+ YDB_LOG_DEBUG_CTX(as, "Updated task cookie",
+ {"taskId", task->GetIdString().data()});
task->Cookie = cookie;
return true;
@@ -481,22 +485,22 @@ TScheduler::TTerminateTaskResult TScheduler::RemoveQueuedTask(ui64 taskId,
{
auto it = Tasks.find(std::make_pair(client, taskId));
if (it == Tasks.end()) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "RemoveQueuedTask failed for task %" PRIu64 " to %s: cannot remove unknown task",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "RemoveQueuedTask failed: cannot remove unknown task",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return TTerminateTaskResult(false, nullptr);
}
auto task = it->second;
if (task->InFly) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "RemoveQueuedTask failed for task %" PRIu64 " to %s: cannot remove in-fly task",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "RemoveQueuedTask failed: cannot remove in-flight task",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return TTerminateTaskResult(false, task);
}
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER, "Removing task %s",
- task->GetIdString().data());
+ YDB_LOG_DEBUG_CTX(as, "Removing queued task",
+ {"taskId", task->GetIdString().data()});
EraseTask(task, false, as);
@@ -510,22 +514,23 @@ TScheduler::TTerminateTaskResult TScheduler::FinishTask(ui64 taskId,
{
auto it = Tasks.find(std::make_pair(client, taskId));
if (it == Tasks.end()) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "FinishTask failed for task %" PRIu64 " to %s: cannot finish unknown task",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "FinishTask failed: cannot finish unknown task",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return TTerminateTaskResult(false, nullptr);
}
auto task = it->second;
if (!task->InFly) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "FinishTask failed for task %" PRIu64 " to %s: cannot finish queued task",
- taskId, ToString(client).c_str());
+ YDB_LOG_DEBUG_CTX(as, "FinishTask failed for task: cannot finish queued task",
+ {"taskId", taskId},
+ {"client", ToString(client)});
return TTerminateTaskResult(false, task);
}
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER, "Finish task %s (release resources {%s})",
- task->GetIdString().data(), JoinSeq(", ", task->RequiredResources).data());
+ YDB_LOG_DEBUG_CTX(as, "Finished task and released resources",
+ {"taskId", task->GetIdString().data()},
+ {"requiredResources", JoinSeq(", ", task->RequiredResources).data()});
// Add execution time to statistics but cancelled tasks
// don't affect average execution time.
@@ -572,14 +577,17 @@ void TScheduler::EraseTask(TTaskPtr task, bool finished, const TActorSystem &as)
queue->EraseTask(task, finished, Now);
if (oldp != queue->PlannedResourceUsage)
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Updated planned resource usage for queue %s from %f to %f (remove task %s)",
- queue->Name.data(), oldp, queue->PlannedResourceUsage, task->GetIdString().data());
+ YDB_LOG_DEBUG_CTX(as, "Updated planned resource usage after removing task",
+ {"queueName", queue->Name.data()},
+ {"previousPlannedResourceUsage", oldp},
+ {"plannedResourceUsage", queue->PlannedResourceUsage},
+ {"taskId", task->GetIdString().data()});
if (oldr != queue->RealResourceUsage)
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Updated real resource usage for queue %s from %f to %f",
- queue->Name.data(), oldr, queue->RealResourceUsage);
+ YDB_LOG_DEBUG_CTX(as, "Updated real resource usage after removing task",
+ {"queueName", queue->Name.data()},
+ {"previousRealResourceUsage", oldr},
+ {"realResourceUsage", queue->RealResourceUsage});
Tasks.erase(std::make_pair(task->Client, task->TaskId));
}
@@ -603,9 +611,8 @@ void TScheduler::ScheduleTasks(const TActorSystem &as,
auto task = queue->FrontTask();
if (task->GetRequiredResourcesMask() & blockedResources) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Skip queue %s blocked by an earlier queue",
- queue->Name.c_str());
+ YDB_LOG_DEBUG_CTX(as, "Skipping queue blocked by an earlier queue",
+ {"queueName", queue->Name});
continue;
}
@@ -614,9 +621,8 @@ void TScheduler::ScheduleTasks(const TActorSystem &as,
// Allow resource over-usage if no tasks are running.
if (!ResourceLimit->HasResources(task->RequiredResources)
&& *TotalCounters->InFlyTasks) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Not enough resources to start task %s",
- task->GetIdString().data());
+ YDB_LOG_DEBUG_CTX(as, "Not enough resources to start task",
+ {"taskId", task->GetIdString().data()});
blockedResources |= task->GetRequiredResourcesMask();
continue;
}
@@ -625,16 +631,15 @@ void TScheduler::ScheduleTasks(const TActorSystem &as,
// Allow resource over-usage if no tasks are running in this queue.
if (!queue->QueueLimit.HasResources(task->RequiredResources)
&& *queue->QueueCounters.InFlyTasks) {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Skip queue %s due to exceeded limits",
- queue->Name.data());
+ YDB_LOG_DEBUG_CTX(as, "Skipping queue due to exceeded limits",
+ {"queueName", queue->Name.data()});
continue;
}
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Allocate resources {%s} for task %s from queue %s",
- JoinSeq(", ", task->RequiredResources).data(),
- task->GetIdString().data(), queue->Name.data());
+ YDB_LOG_DEBUG_CTX(as, "Allocated resources for task",
+ {"requiredResources", JoinSeq(", ", task->RequiredResources).data()},
+ {"taskId", task->GetIdString().data()},
+ {"queueName", queue->Name.data()});
queue->PopTask();
task->InFly = true;
@@ -657,10 +662,11 @@ void TScheduler::UpdateResourceUsage(const TActorSystem &as)
queue->UpdateRealResourceUsage(Now);
if (old != queue->RealResourceUsage)
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Updated real resource usage for queue %s from %f to %f (in-fly consumption {%s})",
- queue->Name.data(), old, queue->RealResourceUsage,
- JoinSeq(", ", queue->QueueLimit.Used).data());
+ YDB_LOG_DEBUG_CTX(as, "Updated real resource usage for in-flight consumption",
+ {"queueName", queue->Name.data()},
+ {"previousRealResourceUsage", old},
+ {"realResourceUsage", queue->RealResourceUsage},
+ {"queueLimitUsed", JoinSeq(", ", queue->QueueLimit.Used).data()});
}
}
@@ -670,13 +676,15 @@ void TScheduler::AssignTask(TTaskPtr &task, const TActorSystem &as)
TTaskQueuePtr queue = TaskConfig(task->Type).Queue;
if (!TaskConfigs.contains(task->Type)) {
- LOG_ERROR(as, NKikimrServices::RESOURCE_BROKER,
- "Assigning %s task '%s' of unknown type '%s' to default queue",
- state.data(), task->GetIdString().data(), task->Type.data());
+ YDB_LOG_ERROR_CTX(as, "Assigning task of unknown type to default queue",
+ {"state", state.data()},
+ {"taskId", task->GetIdString().data()},
+ {"taskType", task->Type.data()});
} else {
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Assigning %s task %s to queue %s",
- state.data(), task->GetIdString().data(), queue->Name.data());
+ YDB_LOG_DEBUG_CTX(as, "Assigning task to queue",
+ {"state", state.data()},
+ {"taskId", task->GetIdString().data()},
+ {"queueName", queue->Name.data()});
}
auto oldp = queue->PlannedResourceUsage;
@@ -701,14 +709,17 @@ void TScheduler::AssignTask(TTaskPtr &task, const TActorSystem &as)
queue->InsertTask(task, Now);
if (oldr != queue->RealResourceUsage)
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Updated real resource usage for queue %s from %f to %f",
- queue->Name.data(), oldr, queue->RealResourceUsage);
+ YDB_LOG_DEBUG_CTX(as, "Updated real resource usage after inserting task",
+ {"queueName", queue->Name.data()},
+ {"previousRealResourceUsage", oldr},
+ {"realResourceUsage", queue->RealResourceUsage});
if (oldp != queue->PlannedResourceUsage)
- LOG_DEBUG(as, NKikimrServices::RESOURCE_BROKER,
- "Updated planned resource usage for queue %s from %f to %f (insert task %s)",
- queue->Name.data(), oldp, queue->PlannedResourceUsage, task->GetIdString().data());
+ YDB_LOG_DEBUG_CTX(as, "Updated planned resource usage after inserting task",
+ {"queueName", queue->Name.data()},
+ {"previousPlannedResourceUsage", oldp},
+ {"plannedResourceUsage", queue->PlannedResourceUsage},
+ {"taskId", task->GetIdString().data()});
}
const TScheduler::TTaskConfig &TScheduler::TaskConfig(const TString &type) const
@@ -964,8 +975,8 @@ bool TResourceBroker::ReduceTaskResourcesInstant(ui64 taskId, const TResourceVal
return true;
}
- LOG_ERROR(*ActorSystem, NKikimrServices::RESOURCE_BROKER,
- "ReduceTaskResourcesInstant failed for task %" PRIu64, taskId);
+ YDB_LOG_ERROR_CTX(*ActorSystem, "ReduceTaskResourcesInstant failed for task",
+ {"taskId", taskId});
return false;
}
}
@@ -1075,9 +1086,9 @@ bool TResourceBroker::FinishTaskInstant(const TEvResourceBroker::TEvFinishTask &
ActorSystem->Send(task.Client, new TEvResourceBroker::TEvResourceAllocated(task.TaskId, task.Cookie));
});
} else {
- LOG_ERROR(*ActorSystem, NKikimrServices::RESOURCE_BROKER,
- "FinishTaskInstant failed for task %" PRIu64 ": %s",
- ev.TaskId, (result.Task ? "cannot finish queued task" : "cannot finish unknown task"));
+ TString error = result.Task ? "cannot finish queued task" : "cannot finish unknown task";
+ YDB_LOG_ERROR_CTX(*ActorSystem, "FinishTaskInstant failed for task: " + error,
+ {"taskId", ev.TaskId});
}
return result.Success;
@@ -1112,7 +1123,7 @@ TResourceBrokerActor::TResourceBrokerActor(const TResourceBrokerConfig &config,
void TResourceBrokerActor::Bootstrap(const TActorContext &ctx)
{
- LOG_DEBUG(ctx, NKikimrServices::RESOURCE_BROKER, "TResourceBrokerActor bootstrap");
+ YDB_LOG_DEBUG_CTX(ctx, "TResourceBrokerActor bootstrap");
NActors::TMon* mon = AppData(ctx)->Mon;
if (mon) {
@@ -1186,13 +1197,15 @@ void TResourceBrokerActor::Handle(TEvResourceBroker::TEvConfigure::TPtr &ev,
{
auto &config = ev->Get()->Record;
if (ev->Get()->Merge) {
- LOG_INFO_S(ctx, NKikimrServices::RESOURCE_BROKER, "New config diff: " << config.ShortDebugString());
+ YDB_LOG_INFO_CTX(ctx, "TEvConfigure: applying config merge",
+ {"configDiff", config.ShortDebugString()});
auto current = ResourceBroker->GetConfig();
MergeConfigUpdates(current, config);
config.Swap(&current);
}
- LOG_INFO_S(ctx, NKikimrServices::RESOURCE_BROKER, "New config: " << config.ShortDebugString());
+ YDB_LOG_INFO_CTX(ctx, "TEvConfigure: applying new config",
+ {"config", config.ShortDebugString()});
TSet<TString> queues;
TSet<TString> tasks;
@@ -1234,10 +1247,8 @@ void TResourceBrokerActor::Handle(TEvResourceBroker::TEvConfigure::TPtr &ev,
ResourceBroker->Configure(std::move(config));
}
- LOG_LOG_S(ctx,
- success ? NActors::NLog::PRI_INFO : NActors::NLog::PRI_ERROR,
- NKikimrServices::RESOURCE_BROKER,
- "Configure result: " << response->Record.ShortDebugString());
+ YDB_LOG_CTX(ctx, success ? NActors::NLog::PRI_INFO : NActors::NLog::PRI_ERROR, "TResourceBrokerActor::Handle TEvConfigure: configuration result",
+ {"result", response->Record.ShortDebugString()});
auto newConfig = ResourceBroker->GetConfig();
for (auto& queue : newConfig.GetQueues()) {