From 44001a42a13e8409f43130a1fe95c88281edbe81 Mon Sep 17 00:00:00 2001 From: Xiaozhen Liu Date: Mon, 10 Aug 2026 17:11:37 -0700 Subject: [PATCH 1/3] feat(amber): carry cache-reuse status as a metrics flag Add a reused_from_cache boolean to OperatorMetrics so a later scheduler PR can report an operator whose results were reused from the operator port cache. Per the direction settled in #5880, reuse is provenance of a completed operator, not a distinct state: a reused operator reports COMPLETED, and this flag carries the distinction to the frontend for display. - executionruntimestate.proto: bool reused_from_cache on OperatorMetrics. - aggregateMetrics: a logical operator is reused only when every one of its physical operators is. - OperatorAggregatedMetrics and the statistics event carry the flag; the TS OperatorStatistics interface gains the matching optional field. Nothing sets the flag until the producer lands with #5884, so with an empty cache the engine behaves identically to before. --- .../engine/common/executionruntimestate.proto | 4 +++ .../execution/ExecutionUtils.scala | 6 +++- .../event/OperatorStatisticsUpdateEvent.scala | 4 ++- .../web/service/ExecutionStatsService.scala | 3 +- .../execution/ExecutionUtilsSpec.scala | 29 +++++++++++++++++-- .../types/execute-workflow.interface.ts | 2 ++ 6 files changed, 43 insertions(+), 5 deletions(-) diff --git a/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto b/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto index e712b3adc8a..e384e80c80e 100644 --- a/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto +++ b/amber/src/main/protobuf/org/apache/texera/amber/engine/common/executionruntimestate.proto @@ -82,6 +82,10 @@ message OperatorStatistics{ message OperatorMetrics{ architecture.rpc.WorkflowAggregatedState operator_state = 1 [(scalapb.field).no_box = true]; OperatorStatistics operator_statistics = 2 [(scalapb.field).no_box = true]; + // True when the operator's results were reused from the operator port cache + // instead of being computed by workers. Provenance of a completed operator, + // not a distinct state; the operator still reports COMPLETED. + bool reused_from_cache = 3; } message ExecutionStatsStore { diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala index 666aeece421..038edc45e07 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala @@ -77,7 +77,11 @@ object ExecutionUtils { dataProcessingTimeSum, controlProcessingTimeSum, idleTimeSum - ) + ), + // A logical operator is reused from cache only when every one of its + // physical operators is. `metrics` is non-empty here, so this cannot + // hold vacuously. + reusedFromCache = metrics.forall(_.reusedFromCache) ) } diff --git a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala index d4aa6117c91..794c687e8e7 100644 --- a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala +++ b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala @@ -30,7 +30,9 @@ case class OperatorAggregatedMetrics( numWorkers: Long, aggregatedDataProcessingTime: Long, aggregatedControlProcessingTime: Long, - aggregatedIdleTime: Long + aggregatedIdleTime: Long, + // Provenance: the operator completed by reusing cached results (no workers ran). + reusedFromCache: Boolean = false ) case class OperatorStatisticsUpdateEvent(operatorStatistics: Map[String, OperatorAggregatedMetrics]) diff --git a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala index f112f4f65da..1242573d110 100644 --- a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala +++ b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala @@ -124,7 +124,8 @@ class ExecutionStatsService( metrics.operatorStatistics.numWorkers, metrics.operatorStatistics.dataProcessingTime, metrics.operatorStatistics.controlProcessingTime, - metrics.operatorStatistics.idleTime + metrics.operatorStatistics.idleTime, + reusedFromCache = metrics.reusedFromCache ) (x._1, res) }) diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala index 237936fa06c..787a466e847 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala @@ -190,11 +190,13 @@ class ExecutionUtilsSpec extends AnyFlatSpec { numWorkers: Int = 0, dataTime: Long = 0, controlTime: Long = 0, - idleTime: Long = 0 + idleTime: Long = 0, + reused: Boolean = false ): OperatorMetrics = OperatorMetrics( state, - OperatorStatistics(input, output, numWorkers, dataTime, controlTime, idleTime) + OperatorStatistics(input, output, numWorkers, dataTime, controlTime, idleTime), + reusedFromCache = reused ) "ExecutionUtils.aggregateMetrics" should "return UNINITIALIZED defaults when given no metrics" in { @@ -337,4 +339,27 @@ class ExecutionUtilsSpec extends AnyFlatSpec { assert(result.operatorStatistics.numWorkers == 3) assert(result.operatorStatistics.dataProcessingTime == 12) } + + // -- aggregateMetrics: reused-from-cache provenance ---------------------- + + it should "report reusedFromCache only when every physical operator is reused" in { + val reusedA = metricsWith(WorkflowAggregatedState.COMPLETED, reused = true) + val reusedB = metricsWith(WorkflowAggregatedState.COMPLETED, reused = true) + val computed = metricsWith(WorkflowAggregatedState.COMPLETED) + + assert(ExecutionUtils.aggregateMetrics(List(reusedA, reusedB)).reusedFromCache) + assert(!ExecutionUtils.aggregateMetrics(List(reusedA, computed)).reusedFromCache) + assert(!ExecutionUtils.aggregateMetrics(List(computed)).reusedFromCache) + } + + it should "default reusedFromCache to false for empty input and untouched metrics" in { + // Empty input takes the early-return path, whose default is false. This is + // the empty-cache property for the flag: nothing sets it until a producer does. + assert(!ExecutionUtils.aggregateMetrics(Iterable.empty).reusedFromCache) + assert( + !ExecutionUtils + .aggregateMetrics(List(metricsWith(WorkflowAggregatedState.RUNNING))) + .reusedFromCache + ) + } } diff --git a/frontend/src/app/workspace/types/execute-workflow.interface.ts b/frontend/src/app/workspace/types/execute-workflow.interface.ts index 8bb7696edf2..54d1a2d6ff6 100644 --- a/frontend/src/app/workspace/types/execute-workflow.interface.ts +++ b/frontend/src/app/workspace/types/execute-workflow.interface.ts @@ -81,6 +81,8 @@ export enum OperatorState { export interface OperatorStatistics extends Readonly<{ operatorState: OperatorState; + // Provenance: the operator completed by reusing cached results (no workers ran). + reusedFromCache?: boolean; aggregatedInputRowCount: number; aggregatedInputSize?: number; inputPortMetrics: Record; From 82f92a028d4f7f51430ed6ba36eaa8066dd4dd8d Mon Sep 17 00:00:00 2001 From: Xiaozhen Liu Date: Wed, 12 Aug 2026 16:40:42 -0700 Subject: [PATCH 2/3] fix(amber): keep reused_from_cache across stats diffing, pin the wire Review follow-ups on #6729: - computeStatsDiff rebuilt OperatorMetrics from its first two fields, an identity transform that silently reset reused_from_cache on the persistence path. Return the combined map as is. - Set the flag non-default in TexeraWebSocketEventSpec's fixture so the symmetric round trip pins it on the wire. - Reword a test comment that relied on PR-history vocabulary and drop a duplicate assertion. --- .../web/service/ExecutionStatsService.scala | 18 +----------------- .../execution/ExecutionUtilsSpec.scala | 11 +++-------- .../event/TexeraWebSocketEventSpec.scala | 5 ++++- 3 files changed, 8 insertions(+), 26 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala index 1242573d110..db10efc4ef5 100644 --- a/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala +++ b/amber/src/main/scala/org/apache/texera/web/service/ExecutionStatsService.scala @@ -237,23 +237,7 @@ class ExecutionStatsService( val updatedLastMetrics = lastPersistedMetrics ++ newKeys.map(_ -> defaultMetrics) // Combine new metrics with old metrics for keys that are no longer present - val completeMetricsMap = newMetrics ++ oldKeys.map(key => key -> updatedLastMetrics(key)) - - // Transform the complete metrics map to ensure consistent structure - completeMetricsMap.map { - case (key, metrics) => - key -> OperatorMetrics( - metrics.operatorState, - OperatorStatistics( - metrics.operatorStatistics.inputMetrics, - metrics.operatorStatistics.outputMetrics, - metrics.operatorStatistics.numWorkers, - metrics.operatorStatistics.dataProcessingTime, - metrics.operatorStatistics.controlProcessingTime, - metrics.operatorStatistics.idleTime - ) - ) - } + newMetrics ++ oldKeys.map(key => key -> updatedLastMetrics(key)) } private def storeRuntimeStatistics( diff --git a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala index 787a466e847..3bdb6c2e860 100644 --- a/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala +++ b/amber/src/test/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtilsSpec.scala @@ -352,14 +352,9 @@ class ExecutionUtilsSpec extends AnyFlatSpec { assert(!ExecutionUtils.aggregateMetrics(List(computed)).reusedFromCache) } - it should "default reusedFromCache to false for empty input and untouched metrics" in { - // Empty input takes the early-return path, whose default is false. This is - // the empty-cache property for the flag: nothing sets it until a producer does. + it should "default reusedFromCache to false for empty input" in { + // Empty input takes the early-return path, whose default is false; metrics + // that no producer has marked keep that default too. assert(!ExecutionUtils.aggregateMetrics(Iterable.empty).reusedFromCache) - assert( - !ExecutionUtils - .aggregateMetrics(List(metricsWith(WorkflowAggregatedState.RUNNING))) - .reusedFromCache - ) } } diff --git a/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala b/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala index 6a1e063d59a..eeeb7f89a9a 100644 --- a/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/model/websocket/event/TexeraWebSocketEventSpec.scala @@ -123,7 +123,10 @@ class TexeraWebSocketEventSpec extends AnyFlatSpec with Matchers { numWorkers = 17L, aggregatedDataProcessingTime = 18L, aggregatedControlProcessingTime = 19L, - aggregatedIdleTime = 20L + aggregatedIdleTime = 20L, + // Non-default on purpose: the symmetric round trip below only pins this + // field on the wire if a drop would change the value read back. + reusedFromCache = true ) private val resultRow = objectMapper.createObjectNode().put("city", "Irvine") From a48b2236a36d9210a05128b8d761d03a20a9c57b Mon Sep 17 00:00:00 2001 From: Xiaozhen Liu Date: Mon, 17 Aug 2026 10:48:16 -0700 Subject: [PATCH 3/3] fix(amber): drop the DTO default, narrow the reuse rollup comment Review follow-ups on #6729: - OperatorAggregatedMetrics.reusedFromCache loses its default so a new construction site must decide the flag explicitly instead of silently sending false, the same shape as the deleted computeStatsDiff rebuild. - The aggregateMetrics comment claimed "every physical operator"; the input holds the operators that currently have a region execution, so it now says so and names the shared transient window. --- .../coordinator/execution/ExecutionUtils.scala | 9 ++++++--- .../websocket/event/OperatorStatisticsUpdateEvent.scala | 4 +++- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala index 038edc45e07..303fba9a1f7 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/coordinator/execution/ExecutionUtils.scala @@ -78,9 +78,12 @@ object ExecutionUtils { controlProcessingTimeSum, idleTimeSum ), - // A logical operator is reused from cache only when every one of its - // physical operators is. `metrics` is non-empty here, so this cannot - // hold vacuously. + // Fully-reused semantics: partial reuse is possible (HashJoin's build and + // probe sit in different regions), and a partially reused operator reports + // false; per-port detail comes from the cache entries. The input holds the + // operators that currently have a region execution, so this shares the + // state field's transient window until all regions exist. Non-empty here, + // so the forall cannot hold vacuously. reusedFromCache = metrics.forall(_.reusedFromCache) ) } diff --git a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala index 794c687e8e7..af91f657843 100644 --- a/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala +++ b/amber/src/main/scala/org/apache/texera/web/model/websocket/event/OperatorStatisticsUpdateEvent.scala @@ -32,7 +32,9 @@ case class OperatorAggregatedMetrics( aggregatedControlProcessingTime: Long, aggregatedIdleTime: Long, // Provenance: the operator completed by reusing cached results (no workers ran). - reusedFromCache: Boolean = false + // Deliberately no default: a new construction site must decide the flag + // explicitly instead of silently sending false. + reusedFromCache: Boolean ) case class OperatorStatisticsUpdateEvent(operatorStatistics: Map[String, OperatorAggregatedMetrics])