This is an automated email from the ASF dual-hosted git repository.

zhoujinsong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/amoro.git


The following commit(s) were added to refs/heads/master by this push:
     new dc126a34c [ams] Fix optimizing status metrics never updated and reset 
on restart (#4356)
dc126a34c is described below

commit dc126a34cfe8950dc63a42a9e4774b9af9297603
Author: ConradJam <[email protected]>
AuthorDate: Mon Sep 7 10:54:56 2026 +0800

    [ams] Fix optimizing status metrics never updated and reset on restart 
(#4356)
    
    The table optimizing status metrics (table_optimizing_status_in_* and
    table_optimizing_status_*_duration_mills) have been broken since the
    TableRuntime storage refactor in AMORO-3747:
    
    1. TableOptimizingMetrics#statusChanged is never invoked, so the status
       gauges stay at their initial values (idle) for the whole process
       lifetime and non-idle durations are always zero.
    2. The runtime is not seeded with the persisted status code and update
       time when constructed, so durations restart from zero after an AMS
       restart.
    3. Same-status status-code writes (e.g. setPendingInput on a non-idle
       table) refresh status_code_update_time, restarting the persisted
       duration clock used by the dashboard.
    
    The fix:
    - seed TableOptimizingMetrics from the persisted status code and
      status_code_update_time when constructing DefaultTableRuntime
    - notify the metrics from the status-change handler after the update
      is committed to the store
    - skip status-code updates that don't change the value, so no-op
      writes no longer refresh status_code_update_time or fire redundant
      status-change notifications
---
 .../amoro/server/table/DefaultTableRuntime.java    | 32 ++++++--
 .../server/table/DefaultTableRuntimeStore.java     | 17 ++++-
 .../amoro/server/table/DefaultTableService.java    |  3 +
 .../table/TestDefaultTableRuntimeHandler.java      | 85 ++++++++++++++++++++++
 4 files changed, 129 insertions(+), 8 deletions(-)

diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntime.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntime.java
index a3467bcd3..fe8813709 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntime.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntime.java
@@ -99,6 +99,13 @@ public class DefaultTableRuntime extends 
AbstractTableRuntime {
     super(store);
     this.optimizingMetrics =
         new TableOptimizingMetrics(store.getTableIdentifier(), 
store.getGroupName());
+    if (store instanceof DefaultTableRuntimeStore) {
+      DefaultTableRuntimeStore runtimeStore = (DefaultTableRuntimeStore) store;
+      OptimizingStatus initialStatus = 
OptimizingStatus.ofCode(runtimeStore.getStatusCode());
+      if (initialStatus != null) {
+        this.optimizingMetrics.statusChanged(initialStatus, 
runtimeStore.getStatusCodeUpdateTime());
+      }
+    }
     this.orphanFilesCleaningMetrics =
         new TableOrphanFilesCleaningMetrics(store.getTableIdentifier());
     this.tableSummaryMetrics = new 
TableSummaryMetrics(store.getTableIdentifier());
@@ -163,6 +170,26 @@ public class DefaultTableRuntime extends 
AbstractTableRuntime {
     return OptimizingStatus.ofCode(getStatusCode());
   }
 
+  /**
+   * Notifies the optimizing metrics after a status-code update has been 
committed to the store.
+   * Invoked from the store's status-change handler callback; no-op when the 
committed update didn't
+   * actually change the status.
+   *
+   * @param originalStatus the status before the committed update.
+   */
+  void onStatusPersisted(OptimizingStatus originalStatus) {
+    OptimizingStatus currentStatus = getOptimizingStatus();
+    if (currentStatus == null || currentStatus == originalStatus) {
+      return;
+    }
+    TableRuntimeStore store = store();
+    long updateTime =
+        store instanceof DefaultTableRuntimeStore
+            ? ((DefaultTableRuntimeStore) store).getStatusCodeUpdateTime()
+            : System.currentTimeMillis();
+    optimizingMetrics.statusChanged(currentStatus, updateTime);
+  }
+
   public long getLastMajorOptimizingTime() {
     return store().getState(OPTIMIZING_STATE_KEY).getLastMajorOptimizingTime();
   }
@@ -319,17 +346,14 @@ public class DefaultTableRuntime extends 
AbstractTableRuntime {
   }
 
   public void beginPlanning() {
-    OptimizingStatus originalStatus = getOptimizingStatus();
     store().begin().updateStatusCode(code -> 
OptimizingStatus.PLANNING.getCode()).commit();
   }
 
   public void planFailed() {
-    OptimizingStatus originalStatus = getOptimizingStatus();
     store().begin().updateStatusCode(code -> 
OptimizingStatus.PENDING.getCode()).commit();
   }
 
   public void beginProcess(OptimizingProcess optimizingProcess) {
-    OptimizingStatus originalStatus = getOptimizingStatus();
     this.optimizingProcess = optimizingProcess;
 
     store()
@@ -343,7 +367,6 @@ public class DefaultTableRuntime extends 
AbstractTableRuntime {
   }
 
   public void completeProcess(boolean success) {
-    OptimizingStatus originalStatus = getOptimizingStatus();
     OptimizingType processType = optimizingProcess.getOptimizingType();
 
     store()
@@ -413,7 +436,6 @@ public class DefaultTableRuntime extends 
AbstractTableRuntime {
   }
 
   public void beginCommitting() {
-    OptimizingStatus originalStatus = getOptimizingStatus();
     store().begin().updateStatusCode(code -> 
OptimizingStatus.COMMITTING.getCode()).commit();
   }
 
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntimeStore.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntimeStore.java
index 3f4d08226..a1354b190 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntimeStore.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableRuntimeStore.java
@@ -41,6 +41,7 @@ import org.slf4j.LoggerFactory;
 
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.locks.Lock;
@@ -102,6 +103,11 @@ public class DefaultTableRuntimeStore extends 
PersistentBase implements TableRun
     return this.meta.getStatusCode();
   }
 
+  /** @return the persisted timestamp when the status code was last updated. */
+  public long getStatusCodeUpdateTime() {
+    return this.meta.getStatusCodeUpdateTime();
+  }
+
   @Override
   public <T> T getState(StateKey<T> key) {
     Preconditions.checkNotNull(key, "TableRuntime state key cannot be null");
@@ -216,11 +222,16 @@ public class DefaultTableRuntimeStore extends 
PersistentBase implements TableRun
     public TableRuntimeOperation updateStatusCode(Function<Integer, Integer> 
updater) {
       operations.add(
           () -> {
-            Integer newStatusCode = updater.apply(oldMeta.getStatusCode());
+            Integer originalStatusCode = oldMeta.getStatusCode();
+            Integer newStatusCode = updater.apply(originalStatusCode);
+            if (Objects.equals(originalStatusCode, newStatusCode)) {
+              return;
+            }
             oldMeta.setStatusCode(newStatusCode);
+            OptimizingStatus originalStatus = 
OptimizingStatus.ofCode(originalStatusCode);
+            handlerCallback.add(
+                handler -> handler.handleTableChanged(tableRuntime, 
originalStatus));
           });
-      OptimizingStatus status = 
OptimizingStatus.ofCode(oldMeta.getStatusCode());
-      handlerCallback.add(handler -> handler.handleTableChanged(tableRuntime, 
status));
       metaOperation = true;
       return this;
     }
diff --git 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableService.java
 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableService.java
index f43b76242..9c9c8a5c4 100644
--- 
a/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableService.java
+++ 
b/amoro-ams/src/main/java/org/apache/amoro/server/table/DefaultTableService.java
@@ -192,6 +192,9 @@ public class DefaultTableService extends PersistentBase 
implements TableService
 
   @Override
   public void handleTableChanged(TableRuntime tableRuntime, OptimizingStatus 
originalStatus) {
+    if (tableRuntime instanceof DefaultTableRuntime) {
+      ((DefaultTableRuntime) tableRuntime).onStatusPersisted(originalStatus);
+    }
     if (headHandler != null) {
       headHandler.fireStatusChanged(tableRuntime, originalStatus);
     }
diff --git 
a/amoro-ams/src/test/java/org/apache/amoro/server/table/TestDefaultTableRuntimeHandler.java
 
b/amoro-ams/src/test/java/org/apache/amoro/server/table/TestDefaultTableRuntimeHandler.java
index 898b3d015..afaabddad 100644
--- 
a/amoro-ams/src/test/java/org/apache/amoro/server/table/TestDefaultTableRuntimeHandler.java
+++ 
b/amoro-ams/src/test/java/org/apache/amoro/server/table/TestDefaultTableRuntimeHandler.java
@@ -30,6 +30,11 @@ import org.apache.amoro.config.Configurations;
 import org.apache.amoro.config.TableConfiguration;
 import org.apache.amoro.hive.catalog.HiveCatalogTestHelper;
 import org.apache.amoro.hive.catalog.HiveTableTestHelper;
+import org.apache.amoro.metrics.Gauge;
+import org.apache.amoro.metrics.Metric;
+import org.apache.amoro.metrics.MetricDefine;
+import org.apache.amoro.metrics.MetricKey;
+import org.apache.amoro.metrics.MetricRegistry;
 import org.apache.amoro.server.manager.EventsManager;
 import org.apache.amoro.server.manager.MetricManager;
 import org.apache.amoro.server.optimizing.OptimizingStatus;
@@ -43,6 +48,7 @@ import org.junit.runner.RunWith;
 import org.junit.runners.Parameterized;
 
 import java.util.List;
+import java.util.Map;
 
 @RunWith(Parameterized.class)
 public class TestDefaultTableRuntimeHandler extends AMSTableTestBase {
@@ -166,6 +172,85 @@ public class TestDefaultTableRuntimeHandler extends 
AMSTableTestBase {
     tableService = null;
   }
 
+  @Test
+  public void testStatusChangeMetricsWiring() throws Exception {
+    tableService = new DefaultTableService(new Configurations(), 
CATALOG_MANAGER, runtimeFactory);
+    tableService.addHandlerChain(new TestHandler());
+    tableService.initialize();
+    if (!(catalogTestHelper().tableFormat().equals(TableFormat.MIXED_HIVE)
+        && TEST_HMS.getHiveClient().getDatabase(TableTestHelper.TEST_DB_NAME) 
!= null)) {
+      createDatabase();
+    }
+    createTable();
+    ServerTableIdentifier tableId = tableManager().listManagedTables().get(0);
+    DefaultTableRuntime runtime = getDefaultTableRuntime(tableId.getId());
+
+    // a brand-new table reports idle in the status gauges
+    Assert.assertEquals(OptimizingStatus.IDLE, runtime.getOptimizingStatus());
+    Assert.assertEquals(1L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_IDLE));
+    Assert.assertEquals(0L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_PLANNING));
+
+    // IDLE -> PLANNING notifies the metrics via the status-change handler 
callback
+    runtime.beginPlanning();
+    Assert.assertEquals(OptimizingStatus.PLANNING, 
runtime.getOptimizingStatus());
+    Assert.assertEquals(1L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_PLANNING));
+    Assert.assertEquals(0L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_IDLE));
+    Thread.sleep(50);
+    Assert.assertTrue(
+        
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_PLANNING_DURATION) > 
0);
+
+    // PLANNING -> PENDING resets the planning duration and moves the pending 
gauges
+    runtime.planFailed();
+    Assert.assertEquals(OptimizingStatus.PENDING, 
runtime.getOptimizingStatus());
+    Assert.assertEquals(1L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_PENDING));
+    Assert.assertEquals(0L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_PLANNING));
+    Assert.assertEquals(
+        0L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_PLANNING_DURATION));
+    Thread.sleep(50);
+    long pendingDurationBeforeRestart =
+        
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_PENDING_DURATION);
+    Assert.assertTrue(pendingDurationBeforeRestart > 0);
+
+    // A same-status write must not restart the persisted duration clock used 
after a restart.
+    DefaultTableRuntimeStore runtimeStore = (DefaultTableRuntimeStore) 
runtime.store();
+    long pendingStatusUpdateTime = runtimeStore.getStatusCodeUpdateTime();
+    runtime.store().begin().updateStatusCode(status -> status).commit();
+    Assert.assertEquals(pendingStatusUpdateTime, 
runtimeStore.getStatusCodeUpdateTime());
+
+    // metrics survive a restart: the restored runtime is seeded from the 
persisted status
+    tableService.dispose();
+    MetricManager.dispose();
+    EventsManager.dispose();
+    tableService = new DefaultTableService(new Configurations(), 
CATALOG_MANAGER, runtimeFactory);
+    tableService.addHandlerChain(new TestHandler());
+    tableService.initialize();
+    DefaultTableRuntime restoredRuntime = 
getDefaultTableRuntime(tableId.getId());
+    Assert.assertEquals(OptimizingStatus.PENDING, 
restoredRuntime.getOptimizingStatus());
+    Assert.assertEquals(1L, 
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_IN_PENDING));
+    Assert.assertTrue(
+        
gaugeValue(TableOptimizingMetrics.TABLE_OPTIMIZING_STATUS_PENDING_DURATION)
+            >= pendingDurationBeforeRestart);
+
+    dropTable();
+    dropDatabase();
+    tableService.dispose();
+    MetricManager.dispose();
+    EventsManager.dispose();
+    tableService = null;
+  }
+
+  private long gaugeValue(MetricDefine define) {
+    MetricRegistry registry = MetricManager.getInstance().getGlobalRegistry();
+    for (Map.Entry<MetricKey, Metric> entry : 
registry.getMetrics().entrySet()) {
+      MetricKey key = entry.getKey();
+      if (define.equals(key.getDefine())
+          && TableTestHelper.TEST_TABLE_NAME.equals(key.valueOfTag("table"))) {
+        return ((Gauge<? extends Number>) 
entry.getValue()).getValue().longValue();
+      }
+    }
+    throw new AssertionError("Gauge not registered for " + define.getName());
+  }
+
   protected DefaultTableService tableService() {
     if (tableService != null) {
       return tableService;

Reply via email to