From c1ae2f615093e25e754b88556f78c8f4bfe5ba24 Mon Sep 17 00:00:00 2001 From: Ruslan Iushchenko Date: Mon, 21 Sep 2026 09:04:41 +0200 Subject: [PATCH 1/5] #796 Add worker status tracking to log statuses of other running tasks on job completion. --- .../pramen/core/pipeline/IngestionJob.scala | 5 + .../orchestrator/OrchestratorImpl.scala | 15 ++- .../core/runner/task/TaskRunnerBase.scala | 39 ++++--- .../runner/task/TaskRunnerMultithreaded.scala | 2 +- .../absa/pramen/core/state/WorkerStatus.scala | 22 ++++ .../core/state/WorkerStatusManager.scala | 37 ++++++ .../za/co/absa/pramen/core/utils/Emoji.scala | 1 + .../core/state/WorkerStatusManagerSuite.scala | 110 ++++++++++++++++++ 8 files changed, 212 insertions(+), 19 deletions(-) create mode 100644 pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatus.scala create mode 100644 pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatusManager.scala create mode 100644 pramen/core/src/test/scala/za/co/absa/pramen/core/state/WorkerStatusManagerSuite.scala diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/pipeline/IngestionJob.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/pipeline/IngestionJob.scala index 14da581df..daa6129b0 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/pipeline/IngestionJob.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/pipeline/IngestionJob.scala @@ -29,6 +29,7 @@ import za.co.absa.pramen.core.metastore.Metastore import za.co.absa.pramen.core.metastore.model.{MetaTable, ReaderMode} import za.co.absa.pramen.core.metastore.peristence.TransientTableManager import za.co.absa.pramen.core.runner.splitter.{ScheduleStrategy, ScheduleStrategySourcing} +import za.co.absa.pramen.core.state.WorkerStatusManager import za.co.absa.pramen.core.utils.ConfigUtils import za.co.absa.pramen.core.utils.Emoji.WARNING import za.co.absa.pramen.core.utils.SparkUtils._ @@ -177,6 +178,7 @@ class IngestionJob(operationDef: OperationDef, conf: Config, jobStarted: Instant, inputRecordCount: Option[Long]): SaveResult = { + WorkerStatusManager.setStatus(s"Running '$name' for '$infoDate'. Writing data...") val stats = metastore.saveTable(outputTable.name, infoDate, df, inputRecordCount) if (!outputTable.format.isRaw) { @@ -190,6 +192,7 @@ class IngestionJob(operationDef: OperationDef, } try { + WorkerStatusManager.setStatus(s"Running '$name' for '$infoDate'. Running post-processing...") source.postProcess( sourceTable.query, outputTable.name, @@ -201,6 +204,7 @@ class IngestionJob(operationDef: OperationDef, case _: AbstractMethodError => log.warn(s"Sources were built using old version of Pramen that does not support post processing. Ignoring...") } + WorkerStatusManager.setStatus(s"Running '$name' for '$infoDate'. Finalizing...") source.close() val jobFinished = Instant.now @@ -233,6 +237,7 @@ class IngestionJob(operationDef: OperationDef, log.info(s"Getting cached record count for '${query.query}' for $from..$to...") getCachedDataFrame(source, query, from, to).count() } else { + WorkerStatusManager.setStatus(s"Running '$name'. Getting record count for for '${query.query}' for $from..$to") source.getRecordCount(sourceTable.query, from, to) } } diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala index aebb946f8..fca60b311 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala @@ -28,7 +28,8 @@ import za.co.absa.pramen.core.pipeline.{Job, JobBase, JobDependency, OperationTy import za.co.absa.pramen.core.runner.jobrunner.ConcurrentJobRunner import za.co.absa.pramen.core.runner.repartitioner.JobRepartitioner import za.co.absa.pramen.core.runner.splitter.ScheduleStrategyUtils.evaluateRunDate -import za.co.absa.pramen.core.state.PipelineState +import za.co.absa.pramen.core.state.{PipelineState, WorkerStatusManager} +import za.co.absa.pramen.core.utils.Emoji import za.co.absa.pramen.core.utils.Emoji._ import java.time.LocalDate @@ -269,6 +270,18 @@ class OrchestratorImpl extends Orchestrator { dependencyResolver.setFailedTable(outputTable.name) } + + if (!isLazy) { + val statuses = WorkerStatusManager.getStatuses + if (statuses.nonEmpty) { + this.synchronized{ + log.info(s"${Emoji.HAMMER_AND_WRENCH} Statuses of other workers:") + statuses.foreach { status => + log.info(s"Thread ${status.threadId}: $status") + } + } + } + } } @throws[ValidationException] diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala index 227c1d2ce..c9ceefcbf 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala @@ -40,7 +40,6 @@ import za.co.absa.pramen.core.pipeline.JobPreRunStatus._ import za.co.absa.pramen.core.pipeline.PipelineDef.{COUNTRY_KEY, ENVIRONMENT_NAME, PIPELINE_NAME_KEY, TENANT_KEY} import za.co.absa.pramen.core.pipeline._ import za.co.absa.pramen.core.runner.splitter.ScheduleStrategyUtils -import za.co.absa.pramen.core.state.PipelineState import za.co.absa.pramen.core.utils.Emoji._ import za.co.absa.pramen.core.utils.SparkUtils._ import za.co.absa.pramen.core.utils.hive.HiveHelper @@ -149,28 +148,34 @@ abstract class TaskRunnerBase(conf: Config, spark.sparkContext.setJobDescription(description) } - task.job.operation.killMaxExecutionTimeSeconds match { - case Some(timeout) if timeout > 0 => - @volatile var runStatus: RunStatus = null + try { + WorkerStatusManager.setStatus(s"Running '${task.job.name}' for '${task.infoDate}'") + + task.job.operation.killMaxExecutionTimeSeconds match { + case Some(timeout) if timeout > 0 => + @volatile var runStatus: RunStatus = null val taskName = task.job.name.replace(' ', '_') val threadName = s"pramen-worker-$taskName-${task.infoDate}" - try { + try { ThreadUtils.runWithTimeout(Duration(timeout, TimeUnit.SECONDS), Duration(sqlCancellationTimeoutSeconds, TimeUnit.SECONDS), threadName = threadName) { - log.info(s"Running ${task.job.name} with the hard timeout = $timeout seconds.") - runStatus = doValidateOrSkipTask(task) + log.info(s"Running '${task.job.name}' with the hard timeout = $timeout seconds.") + runStatus = doValidateOrSkipTask(task) + } + runStatus + } catch { + case NonFatal(ex) => + failTask(task, started, ex) } - runStatus - } catch { - case NonFatal(ex) => - failTask(task, started, ex) - } - case Some(timeout) => - log.error(s"Incorrect timeout for the task: ${task.job.name}. Should be bigger than zero, got: $timeout.") - doValidateOrSkipTask(task) - case None => - doValidateOrSkipTask(task) + case Some(timeout) => + log.error(s"Incorrect timeout for the task: ${task.job.name}. Should be bigger than zero, got: $timeout.") + doValidateOrSkipTask(task) + case None => + doValidateOrSkipTask(task) + } + } finally { + WorkerStatusManager.setFinished() } } diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerMultithreaded.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerMultithreaded.scala index ff3e30de7..866c68dd7 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerMultithreaded.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerMultithreaded.scala @@ -26,7 +26,7 @@ import za.co.absa.pramen.core.bookkeeper.Bookkeeper import za.co.absa.pramen.core.exceptions.FatalErrorWrapper import za.co.absa.pramen.core.journal.Journal import za.co.absa.pramen.core.pipeline.Task -import za.co.absa.pramen.core.state.PipelineState +import za.co.absa.pramen.core.state.{PipelineState, WorkerStatusManager} import za.co.absa.pramen.core.utils.Emoji import java.util.concurrent.Executors.newFixedThreadPool diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatus.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatus.scala new file mode 100644 index 000000000..bd2be12b2 --- /dev/null +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatus.scala @@ -0,0 +1,22 @@ +/* + * Copyright 2022 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package za.co.absa.pramen.core.state + +case class WorkerStatus( + threadId: Long, + status: String + ) diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatusManager.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatusManager.scala new file mode 100644 index 000000000..8550ef47e --- /dev/null +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/state/WorkerStatusManager.scala @@ -0,0 +1,37 @@ +/* + * Copyright 2022 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package za.co.absa.pramen.core.state + +import scala.collection.mutable + +object WorkerStatusManager { + private val workerStatus = new mutable.HashMap[Long, String]() + + def setStatus(status: String): Unit = synchronized { + val threadId = Thread.currentThread().getId + workerStatus(threadId) = status + } + + def setFinished(): Unit = synchronized { + val threadId = Thread.currentThread().getId + workerStatus.remove(threadId) + } + + def getStatuses: Seq[WorkerStatus] = synchronized { + workerStatus.map { case (threadId, status) => WorkerStatus(threadId, status) }.toSeq + } +} diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/Emoji.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/Emoji.scala index 5b8249eae..69fdfd08f 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/Emoji.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/utils/Emoji.scala @@ -31,4 +31,5 @@ object Emoji { val LIGHT_BULB = "\uD83D\uDCA1" val STAR = "\u2B50" val VOLTAGE = "\u26A1" + val HAMMER_AND_WRENCH = "\uD83D\uDEE0\uFE0F" } diff --git a/pramen/core/src/test/scala/za/co/absa/pramen/core/state/WorkerStatusManagerSuite.scala b/pramen/core/src/test/scala/za/co/absa/pramen/core/state/WorkerStatusManagerSuite.scala new file mode 100644 index 000000000..275ae4161 --- /dev/null +++ b/pramen/core/src/test/scala/za/co/absa/pramen/core/state/WorkerStatusManagerSuite.scala @@ -0,0 +1,110 @@ +/* + * Copyright 2022 ABSA Group Limited + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package za.co.absa.pramen.core.state + +import org.scalatest.wordspec.AnyWordSpec + +import java.util.concurrent.{CountDownLatch, TimeUnit} + +class WorkerStatusManagerSuite extends AnyWordSpec { + "setStatus" should { + "register the status of the current thread" in { + try { + WorkerStatusManager.setStatus("running task 1") + + val statuses = WorkerStatusManager.getStatuses + + assert(statuses.nonEmpty) + assert(statuses.exists(s => s.threadId == Thread.currentThread().getId && s.status == "running task 1")) + } finally { + WorkerStatusManager.setFinished() + } + } + + "overwrite the previous status of the same thread" in { + try { + WorkerStatusManager.setStatus("running task 1") + WorkerStatusManager.setStatus("running task 2") + + val statuses = WorkerStatusManager.getStatuses.filter(_.threadId == Thread.currentThread().getId) + + assert(statuses.length == 1) + assert(statuses.head.status == "running task 2") + } finally { + WorkerStatusManager.setFinished() + } + } + } + + "setFinished" should { + "remove the status of the current thread" in { + WorkerStatusManager.setStatus("running task 1") + WorkerStatusManager.setFinished() + + val statuses = WorkerStatusManager.getStatuses.filter(_.threadId == Thread.currentThread().getId) + + assert(statuses.isEmpty) + } + + "do nothing if the current thread has no registered status" in { + WorkerStatusManager.setFinished() + WorkerStatusManager.setFinished() + + val statuses = WorkerStatusManager.getStatuses.filter(_.threadId == Thread.currentThread().getId) + + assert(statuses.isEmpty) + } + } + + "getStatuses" should { + "return an empty collection when no statuses are registered" in { + WorkerStatusManager.setFinished() + + assert(WorkerStatusManager.getStatuses.isEmpty) + } + + "return statuses of all running threads" in { + val threadsStarted = new CountDownLatch(2) + val allowToFinish = new CountDownLatch(1) + + val threads = (1 to 2).map { i => + val thread = new Thread(new Runnable { + override def run(): Unit = { + WorkerStatusManager.setStatus(s"worker $i") + threadsStarted.countDown() + allowToFinish.await(10, TimeUnit.SECONDS) + WorkerStatusManager.setFinished() + } + }) + thread.start() + thread + } + + threadsStarted.await(10, TimeUnit.SECONDS) + + val statuses = WorkerStatusManager.getStatuses + + allowToFinish.countDown() + threads.foreach(_.join(10000)) + + assert(statuses.length == 2) + assert(statuses.map(_.status).sortBy(identity) == Seq("worker 1", "worker 2")) + assert(statuses.map(_.threadId).distinct.length == 2) + assert(WorkerStatusManager.getStatuses.isEmpty) + } + } +} From d0bdbf0238de29cc4c468f63a6ee54c1bab95f00 Mon Sep 17 00:00:00 2001 From: Ruslan Iushchenko Date: Wed, 23 Sep 2026 11:05:04 +0200 Subject: [PATCH 2/5] Fix imports --- .../za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala index c9ceefcbf..6e4d908ee 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala @@ -25,7 +25,7 @@ import za.co.absa.pramen.api.jobdef.Schedule import za.co.absa.pramen.api.lock.TokenLockFactory import za.co.absa.pramen.api.status._ import za.co.absa.pramen.bulkload.BulkLoadStateManager -import za.co.absa.pramen.bulkload.model.{BulkLoadPhase, BulkLoadState} +import za.co.absa.pramen.bulkload.model.BulkLoadPhase import za.co.absa.pramen.core.app.config.RuntimeConfig import za.co.absa.pramen.core.bookkeeper.Bookkeeper import za.co.absa.pramen.core.config.Keys.SQL_QUERY_CANCELLATION_TIMEOUT @@ -40,6 +40,7 @@ import za.co.absa.pramen.core.pipeline.JobPreRunStatus._ import za.co.absa.pramen.core.pipeline.PipelineDef.{COUNTRY_KEY, ENVIRONMENT_NAME, PIPELINE_NAME_KEY, TENANT_KEY} import za.co.absa.pramen.core.pipeline._ import za.co.absa.pramen.core.runner.splitter.ScheduleStrategyUtils +import za.co.absa.pramen.core.state.{PipelineState, WorkerStatusManager} import za.co.absa.pramen.core.utils.Emoji._ import za.co.absa.pramen.core.utils.SparkUtils._ import za.co.absa.pramen.core.utils.hive.HiveHelper From ea0ad92497c65f0edb74bad69a91d0ff946065b6 Mon Sep 17 00:00:00 2001 From: Ruslan Iushchenko Date: Wed, 23 Sep 2026 11:23:58 +0200 Subject: [PATCH 3/5] #796 Mark worker status as finished for tasks running in a child thread (Thanks @coderabbitai!) --- .../pramen/core/runner/task/TaskRunnerBase.scala | 16 ++++++++++------ 1 file changed, 10 insertions(+), 6 deletions(-) diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala index 6e4d908ee..88bc03979 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala @@ -156,23 +156,27 @@ abstract class TaskRunnerBase(conf: Config, case Some(timeout) if timeout > 0 => @volatile var runStatus: RunStatus = null - val taskName = task.job.name.replace(' ', '_') - val threadName = s"pramen-worker-$taskName-${task.infoDate}" + val taskName = task.job.name.replace(' ', '_') + val threadName = s"pramen-worker-$taskName-${task.infoDate}" try { - ThreadUtils.runWithTimeout(Duration(timeout, TimeUnit.SECONDS), Duration(sqlCancellationTimeoutSeconds, TimeUnit.SECONDS), threadName = threadName) { + ThreadUtils.runWithTimeout(Duration(timeout, TimeUnit.SECONDS), Duration(sqlCancellationTimeoutSeconds, TimeUnit.SECONDS), threadName = threadName) { log.info(s"Running '${task.job.name}' with the hard timeout = $timeout seconds.") - runStatus = doValidateOrSkipTask(task) + try { + runStatus = doValidateOrSkipTask(task) + } finally { + WorkerStatusManager.setFinished() + } } runStatus } catch { case NonFatal(ex) => failTask(task, started, ex) } - case Some(timeout) => + case Some(timeout) => log.error(s"Incorrect timeout for the task: ${task.job.name}. Should be bigger than zero, got: $timeout.") doValidateOrSkipTask(task) - case None => + case None => doValidateOrSkipTask(task) } } finally { From ca9230ca1984ce74480df645ab96ece27686970b Mon Sep 17 00:00:00 2001 From: Ruslan Iushchenko Date: Wed, 23 Sep 2026 11:39:23 +0200 Subject: [PATCH 4/5] #796 Move status logging outside the deferred dependency update (Thanks @coderabbitai!) --- .../orchestrator/OrchestratorImpl.scala | 24 +++++++++---------- 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala index fca60b311..f26e75724 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/orchestrator/OrchestratorImpl.scala @@ -121,6 +121,18 @@ class OrchestratorImpl extends Orchestrator { state.addTaskCompletion(taskResults) + if (!isLazy) { + val statuses = WorkerStatusManager.getStatuses + if (statuses.nonEmpty) { + this.synchronized{ + log.info(s"${Emoji.HAMMER_AND_WRENCH} Statuses of other workers:") + statuses.foreach { status => + log.info(s"Thread ${status.threadId}: $status") + } + } + } + } + if (hasFatalErrors || hasCriticalJobFailures) { // In case of a fatal error, we either need to interrupt running threads, or wait for them to return. // In the current implementation we wait for threads to finish, but not start new jobs in running threads. @@ -270,18 +282,6 @@ class OrchestratorImpl extends Orchestrator { dependencyResolver.setFailedTable(outputTable.name) } - - if (!isLazy) { - val statuses = WorkerStatusManager.getStatuses - if (statuses.nonEmpty) { - this.synchronized{ - log.info(s"${Emoji.HAMMER_AND_WRENCH} Statuses of other workers:") - statuses.foreach { status => - log.info(s"Thread ${status.threadId}: $status") - } - } - } - } } @throws[ValidationException] From 5da1b2b17a7caed8dc89f61f28c0cce01eb01c9a Mon Sep 17 00:00:00 2001 From: Ruslan Iushchenko Date: Wed, 23 Sep 2026 11:55:18 +0200 Subject: [PATCH 5/5] #796 Set worker status inside the task execution thread for timeout-bound tasks (Thanks @coderabbitai!) --- .../za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala index 88bc03979..70735e728 100644 --- a/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala +++ b/pramen/core/src/main/scala/za/co/absa/pramen/core/runner/task/TaskRunnerBase.scala @@ -150,8 +150,6 @@ abstract class TaskRunnerBase(conf: Config, } try { - WorkerStatusManager.setStatus(s"Running '${task.job.name}' for '${task.infoDate}'") - task.job.operation.killMaxExecutionTimeSeconds match { case Some(timeout) if timeout > 0 => @volatile var runStatus: RunStatus = null @@ -161,6 +159,7 @@ abstract class TaskRunnerBase(conf: Config, try { ThreadUtils.runWithTimeout(Duration(timeout, TimeUnit.SECONDS), Duration(sqlCancellationTimeoutSeconds, TimeUnit.SECONDS), threadName = threadName) { + WorkerStatusManager.setStatus(s"Running '${task.job.name}' for '${task.infoDate}'") log.info(s"Running '${task.job.name}' with the hard timeout = $timeout seconds.") try { runStatus = doValidateOrSkipTask(task) @@ -175,8 +174,10 @@ abstract class TaskRunnerBase(conf: Config, } case Some(timeout) => log.error(s"Incorrect timeout for the task: ${task.job.name}. Should be bigger than zero, got: $timeout.") + WorkerStatusManager.setStatus(s"Running '${task.job.name}' for '${task.infoDate}'") doValidateOrSkipTask(task) case None => + WorkerStatusManager.setStatus(s"Running '${task.job.name}' for '${task.infoDate}'") doValidateOrSkipTask(task) } } finally {