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

yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new 462361ce28 [#12419] feat(core): add entity change log metrics and 
diagnostic logs (#13388)
462361ce28 is described below

commit 462361ce285f3c578a2397819125f0489bf45212
Author: Qi Yu <[email protected]>
AuthorDate: Tue Sep 22 22:40:39 2026 +0800

    [#12419] feat(core): add entity change log metrics and diagnostic logs 
(#13388)
    
    ### What changes were proposed in this pull request?
    
    - Add entity change log metrics for the database tail, cursor, lag, poll
    age and duration, fetched/delivered/applied records, listener failures,
    and cache recovery.
    - Add correlated diagnostic fields to write, poll, delivery, and cache
    invalidation logs.
    - Document the metrics and current one-time delivery behavior; align the
    server and design documentation.
    
    ### Why are the changes needed?
    
    These signals help operators trace a change across nodes and identify
    where cache invalidation stopped. The current poller has no pending
    delivery or retry state; each listener handles its own recovery.
    
    Fix: #12419
    
    ### Does this PR introduce _any_ user-facing change?
    
    Yes. It adds JMX and Prometheus metrics under the `entity-change-log`
    source and more diagnostic logs. It does not change public APIs or
    configuration keys.
    
    ### How was this patch tested?
    
    - `./gradlew :core:spotlessCheck :core:test --tests
    'org.apache.gravitino.metrics.source.TestEntityChangeLogMetricsSource'
    --tests
    'org.apache.gravitino.storage.relational.TestEntityChangeLogPoller'
    --tests
    'org.apache.gravitino.storage.relational.TestEntityCacheChangeLogListener'
    --tests
    'org.apache.gravitino.storage.relational.TestEntityChangeLogDiagnostics'
    -PskipITs -PskipWeb=true`
    - `git diff --check`
    
    The full `:core:test -PskipITs` run was started locally but stopped
    after the test task produced no progress for several minutes; the
    targeted tests above passed on the latest `apache/main`.
---
 .../source/EntityChangeLogMetricsSource.java       | 123 ++++++++++++++
 .../relational/EntityCacheChangeLogListener.java   |  89 ++++++++--
 .../relational/EntityChangeLogDiagnostics.java     |  48 ++++++
 .../storage/relational/EntityChangeLogPoller.java  |  69 +++++++-
 .../gravitino/storage/relational/JDBCBackend.java  |  10 +-
 .../storage/relational/RelationalEntityStore.java  |  20 ++-
 .../relational/service/ModelMetaService.java       |   5 +
 .../source/TestEntityChangeLogMetricsSource.java   |  85 ++++++++++
 .../TestEntityCacheChangeLogListener.java          |  15 +-
 .../relational/TestEntityChangeLogDiagnostics.java |  70 ++++++++
 .../relational/TestEntityChangeLogPoller.java      | 186 ++++++++++++++++++++-
 ...TestRelationalEntityStoreHierarchicalCache.java |  41 +++++
 design-docs/cache-improvement-design.md            |   6 +
 .../gravitino-entity-cache-multinode-design.md     |  64 +++----
 docs/gravitino-server-config.md                    |  11 +-
 docs/metrics.md                                    |  40 +++++
 16 files changed, 821 insertions(+), 61 deletions(-)

diff --git 
a/core/src/main/java/org/apache/gravitino/metrics/source/EntityChangeLogMetricsSource.java
 
b/core/src/main/java/org/apache/gravitino/metrics/source/EntityChangeLogMetricsSource.java
new file mode 100644
index 0000000000..ffbaed05d7
--- /dev/null
+++ 
b/core/src/main/java/org/apache/gravitino/metrics/source/EntityChangeLogMetricsSource.java
@@ -0,0 +1,123 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you 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 org.apache.gravitino.metrics.source;
+
+import com.codahale.metrics.Counter;
+import com.codahale.metrics.Gauge;
+import com.codahale.metrics.Histogram;
+import com.codahale.metrics.Timer;
+import java.util.concurrent.atomic.AtomicLong;
+
+/** Process-local metrics for the entity change log poller and its 
entity-cache listener. */
+public class EntityChangeLogMetricsSource extends MetricsSource {
+  private final AtomicLong dbTailId = new AtomicLong();
+  private final AtomicLong cursorId = new AtomicLong();
+  private final AtomicLong lastSuccessfulPollMs = new AtomicLong();
+  private final AtomicLong lastSuccessfulTailSampleMs = new AtomicLong();
+  private final Counter pollFailures = getCounter("poll-failures-total");
+  private final Counter tailSampleFailures = 
getCounter("tail-sample-failures-total");
+  private final Counter listenerFailures = 
getCounter("listener-failures-total");
+  private final Counter recordsFetched = getCounter("records-fetched-total");
+  private final Counter recordsDelivered = 
getCounter("records-delivered-total");
+  private final Counter recordsApplied = getCounter("records-applied-total");
+  private final Counter invalidationFailures = 
getCounter("invalidation-failures-total");
+  private final Counter fallbackClears = getCounter("fallback-clears-total");
+  private final Histogram batchSize = getHistogram("batch-size-records");
+  private final Timer pollDuration = getTimer("poll-duration");
+
+  /** Creates and registers the nonblocking gauges for one server's change 
log. */
+  public EntityChangeLogMetricsSource() {
+    super("entity-change-log");
+    registerGauge("db-tail-id", (Gauge<Long>) dbTailId::get);
+    registerGauge("cursor-id", (Gauge<Long>) cursorId::get);
+    registerGauge("record-lag", (Gauge<Long>) () -> Math.max(0, dbTailId.get() 
- cursorId.get()));
+    registerGauge(
+        "seconds-since-last-successful-poll",
+        (Gauge<Long>) () -> secondsSince(lastSuccessfulPollMs.get()));
+    // A failed tail sample keeps the previous db-tail-id, so this gauge is 
what says whether that
+    // value, and record-lag derived from it, still describe the database.
+    registerGauge(
+        "seconds-since-last-successful-tail-sample",
+        (Gauge<Long>) () -> secondsSince(lastSuccessfulTailSampleMs.get()));
+  }
+
+  /** Records the database tail sampled by a poll, without querying from the 
gauge. */
+  public void setDbTailId(long id) {
+    dbTailId.set(id);
+    lastSuccessfulTailSampleMs.set(System.currentTimeMillis());
+  }
+
+  /** Records the cursor after a successful delivery. */
+  public void setCursorId(long id) {
+    cursorId.set(id);
+  }
+
+  /** Records a successful database poll, including an empty result. */
+  public void pollSucceeded(int count) {
+    lastSuccessfulPollMs.set(System.currentTimeMillis());
+    recordsFetched.inc(count);
+    batchSize.update(count);
+  }
+
+  /** Records a failed poll query or cycle. */
+  public void pollFailed() {
+    pollFailures.inc();
+  }
+
+  /** Records a failed database-tail sample while allowing an already fetched 
batch to proceed. */
+  public void tailSampleFailed() {
+    tailSampleFailures.inc();
+  }
+
+  /** Records a listener delivery that failed, attributed by its stable class 
name. */
+  public void listenerFailed(String listenerName) {
+    listenerFailures.inc();
+    getCounter("listener-failures." + listenerName.replace('.', '_') + 
"-total").inc();
+  }
+
+  /** Records the rows delivered to one listener without an exception. */
+  public void recordsDelivered(String listenerName, int count) {
+    recordsDelivered.inc(count);
+    getCounter("records-delivered." + listenerName.replace('.', '_') + 
"-total").inc(count);
+  }
+
+  /** Records targeted entity-cache invalidations that completed successfully. 
*/
+  public void recordsApplied(int count) {
+    recordsApplied.inc(count);
+  }
+
+  /** Records a targeted entity-cache invalidation failure. */
+  public void invalidationFailed() {
+    invalidationFailures.inc();
+  }
+
+  /** Records a successful full-cache clear used as a recovery fallback. */
+  public void fallbackCleared() {
+    fallbackClears.inc();
+  }
+
+  /** Starts a poll duration measurement. */
+  public Timer.Context timePoll() {
+    return pollDuration.time();
+  }
+
+  private static long secondsSince(long epochMs) {
+    return epochMs == 0 ? -1 : Math.max(0, (System.currentTimeMillis() - 
epochMs) / 1000);
+  }
+}
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/EntityCacheChangeLogListener.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityCacheChangeLogListener.java
index ad9d568f10..d208d2899a 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/EntityCacheChangeLogListener.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityCacheChangeLogListener.java
@@ -21,9 +21,11 @@ package org.apache.gravitino.storage.relational;
 import com.google.common.base.Preconditions;
 import java.util.List;
 import java.util.Locale;
+import java.util.concurrent.TimeUnit;
 import org.apache.gravitino.Entity.EntityType;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.cache.EntityCache;
+import org.apache.gravitino.metrics.source.EntityChangeLogMetricsSource;
 import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -80,24 +82,50 @@ public class EntityCacheChangeLogListener implements 
EntityChangeLogListener {
   }
 
   private final Target target;
+  private final EntityChangeLogMetricsSource metrics;
 
   /**
-   * Creates a listener that invalidates the given entity store cache directly.
+   * Creates a listener that invalidates the given entity store cache 
directly. Metrics from this
+   * constructor are local to the listener and are not exported by the server 
metrics system.
    *
    * @param cache the per-node entity store cache to keep coherent
    */
   public EntityCacheChangeLogListener(EntityCache cache) {
-    this(asTarget(cache));
+    this(asTarget(cache), new EntityChangeLogMetricsSource());
   }
 
   /**
-   * Creates a listener that invalidates through the given target.
+   * Creates a listener that invalidates the given entity store cache 
directly, with metrics shared
+   * with the poller.
+   *
+   * @param cache the per-node entity store cache to keep coherent
+   * @param metrics process-local change log metrics
+   */
+  public EntityCacheChangeLogListener(EntityCache cache, 
EntityChangeLogMetricsSource metrics) {
+    this(asTarget(cache), metrics);
+  }
+
+  /**
+   * Creates a listener that invalidates through the given target. Metrics 
from this constructor are
+   * local to the listener and are not exported by the server metrics system.
    *
    * @param target the invalidation entry points of the per-node cache to keep 
coherent
    */
   public EntityCacheChangeLogListener(Target target) {
+    this(target, new EntityChangeLogMetricsSource());
+  }
+
+  /**
+   * Creates a listener that invalidates through the given target, with 
metrics shared with the
+   * poller.
+   *
+   * @param target the invalidation entry points of the per-node cache to keep 
coherent
+   * @param metrics process-local change log metrics
+   */
+  public EntityCacheChangeLogListener(Target target, 
EntityChangeLogMetricsSource metrics) {
     Preconditions.checkArgument(target != null, "target cannot be null");
     this.target = target;
+    this.metrics = Preconditions.checkNotNull(metrics, "metrics cannot be 
null");
   }
 
   private static Target asTarget(EntityCache cache) {
@@ -117,43 +145,75 @@ public class EntityCacheChangeLogListener implements 
EntityChangeLogListener {
 
   @Override
   public void onEntityChange(List<EntityChangeRecord> changes) {
+    long startNanos = System.nanoTime();
+    int applied = 0;
+    int skipped = 0;
     for (EntityChangeRecord change : changes) {
       EntityType type = entityType(change);
       NameIdentifier ident = identifier(change);
       if (type == null || ident == null) {
         // Already logged by the parsing helpers. A row that names no entity 
cannot invalidate
         // anything, so skipping it leaves no stale entry behind.
+        skipped++;
         continue;
       }
 
       try {
-        LOG.debug("Invalidating entity cache due to entity change log: {} 
({})", ident, type);
+        LOG.debug(
+            "entityChangeLog invalidate changeId={} entityType={} 
operateType={} ident={} fullName={}",
+            change.getId(),
+            type,
+            change.getOperateType(),
+            ident,
+            change.getFullName());
         target.invalidate(ident, type);
+        applied++;
+        metrics.recordsApplied(1);
       } catch (RuntimeException e) {
+        metrics.invalidationFailed();
         // Dropping a single invalidation would leave this node serving that 
entity stale until it
         // expires. Clearing the whole cache is the safe superset, and it also 
covers the rest of
         // this batch, so there is nothing left to replay.
         LOG.error(
-            "Failed to invalidate {} ({}) from the entity change log, clearing 
the local entity "
-                + "cache to stay coherent",
-            ident,
+            "entityChangeLog targeted invalidation failed changeId={} 
entityType={} "
+                + "operateType={} ident={} fullName={}; clearing full local 
entity cache",
+            change.getId(),
             type,
+            change.getOperateType(),
+            ident,
+            change.getFullName(),
             e);
         target.clear();
+        metrics.fallbackCleared();
+        LOG.debug(
+            "entityChangeLog invalidate batch count={} applied={} skipped={} 
fallbackClear=true durationMs={}",
+            changes.size(),
+            applied,
+            skipped,
+            TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos));
         return;
       }
     }
+    LOG.debug(
+        "entityChangeLog invalidate batch count={} applied={} skipped={} 
fallbackClear=false durationMs={}",
+        changes.size(),
+        applied,
+        skipped,
+        TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNanos));
   }
 
   private EntityType entityType(EntityChangeRecord change) {
     if (change.getEntityType() == null) {
-      LOG.warn("Invalid entity type in entity change log: null");
+      LOG.warn("entityChangeLog malformed changeId={} field=entityType 
value=null", change.getId());
       return null;
     }
     try {
       return 
EntityType.valueOf(change.getEntityType().toUpperCase(Locale.ROOT));
     } catch (IllegalArgumentException e) {
-      LOG.warn("Unknown entity type in entity change log: {}", 
change.getEntityType());
+      LOG.warn(
+          "entityChangeLog malformed changeId={} field=entityType value={}",
+          change.getId(),
+          change.getEntityType());
       return null;
     }
   }
@@ -161,13 +221,20 @@ public class EntityCacheChangeLogListener implements 
EntityChangeLogListener {
   private NameIdentifier identifier(EntityChangeRecord change) {
     String fullName = change.getFullName();
     if (fullName == null || fullName.isEmpty()) {
-      LOG.warn("Invalid full name in entity change log: {}", fullName);
+      LOG.warn(
+          "entityChangeLog malformed changeId={} field=fullName value={}",
+          change.getId(),
+          fullName);
       return null;
     }
     try {
       return EntityChangeLogNameIdentifierCodec.decode(fullName);
     } catch (IllegalArgumentException e) {
-      LOG.warn("Undecodable full name in entity change log: {}", fullName, e);
+      LOG.warn(
+          "entityChangeLog malformed changeId={} field=fullName value={}",
+          change.getId(),
+          fullName,
+          e);
       return null;
     }
   }
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogDiagnostics.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogDiagnostics.java
new file mode 100644
index 0000000000..00b2334e84
--- /dev/null
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogDiagnostics.java
@@ -0,0 +1,48 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you 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 org.apache.gravitino.storage.relational;
+
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/** Diagnostic logging for entity change rows appended to the current 
transaction. */
+public final class EntityChangeLogDiagnostics {
+  private static final Logger LOG = 
LoggerFactory.getLogger(EntityChangeLogDiagnostics.class);
+
+  private EntityChangeLogDiagnostics() {}
+
+  /**
+   * Logs a successful append without implying that the enclosing transaction 
committed.
+   *
+   * @param metalake the metalake name
+   * @param entityType the entity type
+   * @param operateType the change operation
+   * @param fullName the encoded identifier stored in the row
+   */
+  public static void logAppended(
+      String metalake, String entityType, OperateType operateType, String 
fullName) {
+    LOG.debug(
+        "entityChangeLog appendedToTransaction metalake={} entityType={} 
operateType={} fullName={}",
+        metalake,
+        entityType,
+        operateType,
+        fullName);
+  }
+}
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java
index 8fa5f51a60..b42b57dc1a 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java
@@ -18,6 +18,7 @@
  */
 package org.apache.gravitino.storage.relational;
 
+import com.codahale.metrics.Timer;
 import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Preconditions;
 import java.util.List;
@@ -26,6 +27,7 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import javax.annotation.Nullable;
+import org.apache.gravitino.metrics.source.EntityChangeLogMetricsSource;
 import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
 import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
 import org.apache.gravitino.storage.relational.utils.SessionUtils;
@@ -75,17 +77,30 @@ public class EntityChangeLogPoller implements AutoCloseable 
{
 
   private final List<EntityChangeLogListener> listeners = new 
CopyOnWriteArrayList<>();
   private final long pollIntervalSecs;
+  private final EntityChangeLogMetricsSource metrics;
 
   private ScheduledExecutorService scheduler;
   private volatile long entityPollHighWaterId = 0;
 
   /**
-   * Creates an {@link EntityChangeLogPoller}.
+   * Creates an {@link EntityChangeLogPoller} with an unregistered metrics 
source for callers that
+   * do not use the server metrics system.
    *
    * @param pollIntervalSecs interval between successive polling cycles
    */
   public EntityChangeLogPoller(long pollIntervalSecs) {
+    this(pollIntervalSecs, new EntityChangeLogMetricsSource());
+  }
+
+  /**
+   * Creates a poller using the metrics source registered by the entity store.
+   *
+   * @param pollIntervalSecs interval between successive polling cycles
+   * @param metrics process-local change log metrics
+   */
+  public EntityChangeLogPoller(long pollIntervalSecs, 
EntityChangeLogMetricsSource metrics) {
     Preconditions.checkArgument(pollIntervalSecs > 0, "pollIntervalSecs must 
be positive");
+    this.metrics = Preconditions.checkNotNull(metrics, "metrics cannot be 
null");
     this.pollIntervalSecs = pollIntervalSecs;
   }
 
@@ -134,6 +149,8 @@ public class EntityChangeLogPoller implements AutoCloseable 
{
         getOrDefault(
             SessionUtils.getWithoutCommit(
                 EntityChangeLogMapper.class, 
EntityChangeLogMapper::selectMaxChangeId));
+    metrics.setDbTailId(entityPollHighWaterId);
+    metrics.setCursorId(entityPollHighWaterId);
     LOG.info(
         "Starting entity change log poller at high-water id {} with a {} 
second interval, "
             + "{} listener(s) registered",
@@ -174,7 +191,7 @@ public class EntityChangeLogPoller implements AutoCloseable 
{
 
   @VisibleForTesting
   void pollChanges() {
-    try {
+    try (Timer.Context ignored = metrics.timePoll()) {
       doPollChanges();
     } catch (Throwable e) {
       // Catch Throwable, not Exception: this method is the task handed to
@@ -185,6 +202,7 @@ public class EntityChangeLogPoller implements AutoCloseable 
{
       if (handleInterruptIfAny(e, "Entity change poll")) {
         return;
       }
+      metrics.pollFailed();
       LOG.warn("Entity change poll failed at high-water id {}", 
entityPollHighWaterId, e);
     }
   }
@@ -198,8 +216,32 @@ public class EntityChangeLogPoller implements 
AutoCloseable {
 
   @Nullable
   private BatchDelivery fetchNextDelivery() {
+    long fetchStartNanos = System.nanoTime();
     List<EntityChangeRecord> changes = fetchEntityChanges();
+    // The tail is for observability only. A failed sample must not suppress 
delivery of rows
+    // already fetched successfully or hold the cursor back.
+    @Nullable Long dbTailId = null;
+    try {
+      dbTailId =
+          getOrDefault(
+              SessionUtils.getWithoutCommit(
+                  EntityChangeLogMapper.class, 
EntityChangeLogMapper::selectMaxChangeId));
+      metrics.setDbTailId(dbTailId);
+    } catch (RuntimeException e) {
+      if (handleInterruptIfAny(e, "Entity change log tail sample")) {
+        throw e;
+      }
+      metrics.tailSampleFailed();
+      LOG.warn("Could not sample entity change log tail; retaining the 
previous gauge value", e);
+    }
+    metrics.pollSucceeded(changes.size());
+    long durationMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - 
fetchStartNanos);
     if (changes.isEmpty()) {
+      LOG.debug(
+          "entityChangeLog poll cursor={} fetched=0 tailId={} durationMs={}",
+          entityPollHighWaterId,
+          dbTailId,
+          durationMs);
       return null;
     }
 
@@ -208,11 +250,13 @@ public class EntityChangeLogPoller implements 
AutoCloseable {
     BatchDelivery delivery =
         new BatchDelivery(immutableChanges, lastChangeId, 
List.copyOf(listeners));
     LOG.debug(
-        "Fetched {} entity change log record(s) after cursor {}, id range [{}, 
{}]: {}",
-        immutableChanges.size(),
+        "entityChangeLog poll cursor={} fetched={} firstId={} lastId={} 
tailId={} durationMs={} records={}",
         entityPollHighWaterId,
+        immutableChanges.size(),
         delivery.firstChangeId(),
         delivery.lastChangeId,
+        dbTailId,
+        durationMs,
         summarize(immutableChanges));
     return delivery;
   }
@@ -280,6 +324,7 @@ public class EntityChangeLogPoller implements AutoCloseable 
{
   private void advanceCursor(BatchDelivery delivery) {
     long previousHighWaterId = entityPollHighWaterId;
     entityPollHighWaterId = delivery.lastChangeId;
+    metrics.setCursorId(entityPollHighWaterId);
     LOG.info(
         "Consumed {} entity change log record(s), id range [{}, {}]; cursor 
advanced from {} to {}; "
             + "newest record is ~{} ms old",
@@ -308,13 +353,21 @@ public class EntityChangeLogPoller implements 
AutoCloseable {
       }
 
       try {
+        LOG.debug(
+            "entityChangeLog delivery listener={} firstId={} lastId={} 
count={} attempt=1",
+            listener.getClass().getName(),
+            delivery.firstChangeId(),
+            delivery.lastChangeId,
+            delivery.changes.size());
         listener.onEntityChange(delivery.changes);
+        metrics.recordsDelivered(listenerMetricName(listener), 
delivery.changes.size());
         LOG.debug(
             "Entity change log listener {} consumed batch id range [{}, {}]",
             listener.getClass().getName(),
             delivery.firstChangeId(),
             delivery.lastChangeId);
       } catch (Throwable e) {
+        metrics.listenerFailed(listenerMetricName(listener));
         // Throwable, not Exception: one faulty listener must not take down 
the whole poller, even
         // if it fails with an Error rather than an Exception.
         LOG.error(
@@ -328,6 +381,14 @@ public class EntityChangeLogPoller implements 
AutoCloseable {
     }
   }
 
+  /** Uses a bounded, stable metric bucket for lambda and anonymous listener 
implementations. */
+  private static String listenerMetricName(EntityChangeLogListener listener) {
+    Class<?> listenerClass = listener.getClass();
+    return listenerClass.isSynthetic() || listenerClass.isAnonymousClass()
+        ? "anonymous"
+        : listenerClass.getName();
+  }
+
   private static class BatchDelivery {
     private final List<EntityChangeRecord> changes;
     private final long lastChangeId;
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java 
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
index c847c8f910..291f00535d 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/JDBCBackend.java
@@ -1099,14 +1099,12 @@ public class JDBCBackend implements RelationalBackend, 
SupportsOrphanedRelationC
 
   private static void insertEntityChange(
       NameIdentifier ident, Entity.EntityType entityType, OperateType 
operateType) {
+    String metalake = NameIdentifierUtil.getMetalake(ident);
+    String fullName = EntityChangeLogNameIdentifierCodec.encode(ident);
     SessionUtils.doWithoutCommit(
         EntityChangeLogMapper.class,
-        mapper ->
-            mapper.insertEntityChange(
-                NameIdentifierUtil.getMetalake(ident),
-                entityType.name(),
-                EntityChangeLogNameIdentifierCodec.encode(ident),
-                operateType));
+        mapper -> mapper.insertEntityChange(metalake, entityType.name(), 
fullName, operateType));
+    EntityChangeLogDiagnostics.logAppended(metalake, entityType.name(), 
operateType, fullName);
   }
 
   private static boolean shouldRecordEntityDrop(Entity.EntityType entityType) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java
index 186d0b8cc8..eda8be4d96 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/RelationalEntityStore.java
@@ -39,6 +39,7 @@ import org.apache.gravitino.Configs;
 import org.apache.gravitino.Entity;
 import org.apache.gravitino.EntityAlreadyExistsException;
 import org.apache.gravitino.EntityStore;
+import org.apache.gravitino.GravitinoEnv;
 import org.apache.gravitino.HasIdentifier;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
@@ -55,6 +56,8 @@ import org.apache.gravitino.cache.EntityCache;
 import org.apache.gravitino.cache.EntityCacheKey;
 import org.apache.gravitino.cache.NoOpsCache;
 import org.apache.gravitino.exceptions.NoSuchEntityException;
+import org.apache.gravitino.metrics.MetricsSystem;
+import org.apache.gravitino.metrics.source.EntityChangeLogMetricsSource;
 import org.apache.gravitino.storage.relational.service.EntityIdService;
 import org.apache.gravitino.utils.Executable;
 import org.slf4j.Logger;
@@ -76,6 +79,9 @@ public class RelationalEntityStore
   private EntityChangeLogPoller entityChangeLogPoller;
   private EntityChangeLogCleaner entityChangeLogCleaner;
   private EntityCache cache;
+  // Created with the store rather than in initialize(), so that a listener 
built before or without
+  // initialization still has somewhere to record. initialize() only registers 
it for export.
+  private final EntityChangeLogMetricsSource changeLogMetrics = new 
EntityChangeLogMetricsSource();
 
   // Advanced before every invalidation observed by this store, whether local 
or replayed from the
   // change log. A shared cache without a local change-log listener needs its 
own distributed
@@ -108,8 +114,13 @@ public class RelationalEntityStore
 
     // Polling and cleanup use separate single-threaded schedulers. Polling 
only dispatches changes
     // to local listeners, while cleanup independently removes records beyond 
the retention period.
+    MetricsSystem metricsSystem = GravitinoEnv.getInstance().metricsSystem();
+    if (metricsSystem != null) {
+      metricsSystem.register(changeLogMetrics);
+    }
     this.entityChangeLogPoller =
-        new 
EntityChangeLogPoller(config.get(Configs.ENTITY_CHANGE_LOG_POLL_INTERVAL_SECS));
+        new EntityChangeLogPoller(
+            config.get(Configs.ENTITY_CHANGE_LOG_POLL_INTERVAL_SECS), 
changeLogMetrics);
     this.entityChangeLogCleaner =
         new EntityChangeLogCleaner(
             
TimeUnit.SECONDS.toMillis(config.get(Configs.ENTITY_CHANGE_LOG_RETENTION_SECS)),
@@ -157,7 +168,8 @@ public class RelationalEntityStore
           public void clear() {
             clearCache();
           }
-        });
+        },
+        changeLogMetrics);
   }
 
   private RelationalBackend createRelationalEntityBackend(Config config) {
@@ -338,6 +350,10 @@ public class RelationalEntityStore
     failure = closeComponent(failure, "entity change log cleaner", 
entityChangeLogCleaner);
     failure = closeComponent(failure, "relational garbage collector", 
garbageCollector);
     failure = closeComponent(failure, "relational backend", backend);
+    MetricsSystem metricsSystem = GravitinoEnv.getInstance().metricsSystem();
+    if (metricsSystem != null) {
+      metricsSystem.unregister(changeLogMetrics);
+    }
 
     if (failure != null) {
       throw failure;
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
index 386ff0f9fc..49f5e7d260 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
@@ -38,6 +38,7 @@ import org.apache.gravitino.exceptions.NoSuchEntityException;
 import org.apache.gravitino.meta.ModelEntity;
 import org.apache.gravitino.meta.NamespacedEntityId;
 import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.EntityChangeLogDiagnostics;
 import 
org.apache.gravitino.storage.relational.EntityChangeLogNameIdentifierCodec;
 import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
 import org.apache.gravitino.storage.relational.mapper.ModelMetaMapper;
@@ -150,6 +151,8 @@ public class ModelMetaService {
                         Entity.EntityType.MODEL.name(),
                         modelFullName,
                         OperateType.DROP));
+            EntityChangeLogDiagnostics.logAppended(
+                metalakeName, Entity.EntityType.MODEL.name(), 
OperateType.DROP, modelFullName);
           });
     } catch (NoSuchEntityException e) {
       // Another writer dropped the model between the read above and this 
transaction. A drop that
@@ -373,6 +376,8 @@ public class ModelMetaService {
                               Entity.EntityType.MODEL.name(),
                               oldFullName,
                               OperateType.ALTER));
+                  EntityChangeLogDiagnostics.logAppended(
+                      metalakeName, Entity.EntityType.MODEL.name(), 
OperateType.ALTER, oldFullName);
                 }
               });
     } catch (RuntimeException re) {
diff --git 
a/core/src/test/java/org/apache/gravitino/metrics/source/TestEntityChangeLogMetricsSource.java
 
b/core/src/test/java/org/apache/gravitino/metrics/source/TestEntityChangeLogMetricsSource.java
new file mode 100644
index 0000000000..98b9fdc6f2
--- /dev/null
+++ 
b/core/src/test/java/org/apache/gravitino/metrics/source/TestEntityChangeLogMetricsSource.java
@@ -0,0 +1,85 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you 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 org.apache.gravitino.metrics.source;
+
+import com.codahale.metrics.Gauge;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+/** Tests the nonblocking state gauges and counters of the entity change log 
source. */
+public class TestEntityChangeLogMetricsSource {
+  @Test
+  void testPollAndFailureMetrics() {
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    Assertions.assertEquals(-1L, gauge(metrics, 
"seconds-since-last-successful-poll"));
+    Assertions.assertEquals(-1L, gauge(metrics, 
"seconds-since-last-successful-tail-sample"));
+
+    metrics.setCursorId(5);
+    metrics.setDbTailId(8);
+    metrics.pollSucceeded(3);
+    metrics.recordsDelivered("org.apache.gravitino.CacheListener", 6);
+    metrics.recordsApplied(6);
+    metrics.pollFailed();
+    metrics.tailSampleFailed();
+    metrics.listenerFailed("org.apache.gravitino.CacheListener");
+    metrics.invalidationFailed();
+    metrics.fallbackCleared();
+
+    Assertions.assertEquals(5L, gauge(metrics, "cursor-id"));
+    Assertions.assertEquals(8L, gauge(metrics, "db-tail-id"));
+    Assertions.assertEquals(3L, gauge(metrics, "record-lag"));
+    Assertions.assertTrue(gauge(metrics, "seconds-since-last-successful-poll") 
>= 0);
+    Assertions.assertTrue(gauge(metrics, 
"seconds-since-last-successful-tail-sample") >= 0);
+    Assertions.assertEquals(
+        3, 
metrics.getMetricRegistry().counter("records-fetched-total").getCount());
+    Assertions.assertEquals(
+        6, 
metrics.getMetricRegistry().counter("records-delivered-total").getCount());
+    Assertions.assertEquals(
+        6,
+        metrics
+            .getMetricRegistry()
+            
.counter("records-delivered.org_apache_gravitino_CacheListener-total")
+            .getCount());
+    Assertions.assertEquals(
+        6, 
metrics.getMetricRegistry().counter("records-applied-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("poll-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("tail-sample-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("listener-failures-total").getCount());
+    Assertions.assertEquals(
+        1,
+        metrics
+            .getMetricRegistry()
+            
.counter("listener-failures.org_apache_gravitino_CacheListener-total")
+            .getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("invalidation-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("fallback-clears-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().histogram("batch-size-records").getCount());
+  }
+
+  private static long gauge(EntityChangeLogMetricsSource metrics, String name) 
{
+    Gauge<?> gauge = metrics.getMetricRegistry().getGauges().get(name);
+    return ((Number) gauge.getValue()).longValue();
+  }
+}
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityCacheChangeLogListener.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityCacheChangeLogListener.java
index 7e8a4e07e5..7d86790171 100644
--- 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityCacheChangeLogListener.java
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityCacheChangeLogListener.java
@@ -38,6 +38,7 @@ import org.apache.gravitino.meta.CatalogEntity;
 import org.apache.gravitino.meta.SchemaEntity;
 import org.apache.gravitino.meta.TableEntity;
 import org.apache.gravitino.meta.TagEntity;
+import org.apache.gravitino.metrics.source.EntityChangeLogMetricsSource;
 import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
 import org.apache.gravitino.storage.relational.po.cache.OperateType;
 import org.apache.gravitino.utils.TestUtil;
@@ -63,12 +64,15 @@ public class TestEntityCacheChangeLogListener {
     cache.put(catalog);
     cache.put(schema);
 
-    EntityCacheChangeLogListener listener = new 
EntityCacheChangeLogListener(cache);
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    EntityCacheChangeLogListener listener = new 
EntityCacheChangeLogListener(cache, metrics);
     listener.onEntityChange(
         List.of(record(EntityType.SCHEMA, schema.nameIdentifier().toString(), 
OperateType.DROP)));
 
     Assertions.assertFalse(cache.contains(schema.nameIdentifier(), 
EntityType.SCHEMA));
     Assertions.assertTrue(cache.contains(catalog.nameIdentifier(), 
EntityType.CATALOG));
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("records-applied-total").getCount());
   }
 
   @Test
@@ -192,7 +196,8 @@ public class TestEntityCacheChangeLogListener {
     NameIdentifier failing = NameIdentifier.of("m1", "boom");
     doThrow(new RuntimeException("boom")).when(cache).invalidate(failing, 
EntityType.CATALOG);
 
-    EntityCacheChangeLogListener listener = new 
EntityCacheChangeLogListener(cache);
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    EntityCacheChangeLogListener listener = new 
EntityCacheChangeLogListener(cache, metrics);
     listener.onEntityChange(
         List.of(
             record(EntityType.CATALOG, "m1.boom", OperateType.DROP),
@@ -202,6 +207,12 @@ public class TestEntityCacheChangeLogListener {
     // replayed and no entry can survive stale.
     verify(cache).clear();
     verify(cache, never()).invalidate(NameIdentifier.of("m1", "ok"), 
EntityType.CATALOG);
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("invalidation-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("fallback-clears-total").getCount());
+    Assertions.assertEquals(
+        0, 
metrics.getMetricRegistry().counter("records-applied-total").getCount());
   }
 
   @Test
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogDiagnostics.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogDiagnostics.java
new file mode 100644
index 0000000000..a5520530b4
--- /dev/null
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogDiagnostics.java
@@ -0,0 +1,70 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you 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 org.apache.gravitino.storage.relational;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.gravitino.storage.relational.po.cache.OperateType;
+import org.apache.logging.log4j.Level;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.core.LogEvent;
+import org.apache.logging.log4j.core.LoggerContext;
+import org.apache.logging.log4j.core.appender.AbstractAppender;
+import org.apache.logging.log4j.core.config.Configuration;
+import org.apache.logging.log4j.core.config.LoggerConfig;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+/** Tests the diagnostic fields without depending on a fully formatted log 
line. */
+public class TestEntityChangeLogDiagnostics {
+  @Test
+  void testAppendLogCarriesChangeIdentityAndTransactionState() {
+    LoggerContext context =
+        (LoggerContext)
+            
LogManager.getContext(EntityChangeLogDiagnostics.class.getClassLoader(), false);
+    Configuration config = context.getConfiguration();
+    List<LogEvent> events = new ArrayList<>();
+    AbstractAppender appender =
+        new AbstractAppender("entityChangeLogCapture", null, null, true, null) 
{
+          @Override
+          public void append(LogEvent event) {
+            events.add(event.toImmutable());
+          }
+        };
+    appender.start();
+    LoggerConfig logger =
+        new LoggerConfig(EntityChangeLogDiagnostics.class.getName(), 
Level.DEBUG, false);
+    logger.addAppender(appender, Level.DEBUG, null);
+    config.addLogger(EntityChangeLogDiagnostics.class.getName(), logger);
+    context.updateLoggers();
+    try {
+      EntityChangeLogDiagnostics.logAppended("ml", "TABLE", OperateType.ALTER, 
"encoded-name");
+      Assertions.assertEquals(1, events.size());
+      Assertions.assertTrue(
+          
events.get(0).getMessage().getFormattedMessage().contains("appendedToTransaction"));
+      Assertions.assertArrayEquals(
+          new Object[] {"ml", "TABLE", OperateType.ALTER, "encoded-name"},
+          events.get(0).getMessage().getParameters());
+    } finally {
+      config.removeLogger(EntityChangeLogDiagnostics.class.getName());
+      context.updateLoggers();
+      appender.stop();
+    }
+  }
+}
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogPoller.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogPoller.java
index 37c05bd18b..2d68d14d68 100644
--- 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogPoller.java
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogPoller.java
@@ -28,6 +28,7 @@ import java.util.ArrayList;
 import java.util.List;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Function;
+import org.apache.gravitino.metrics.source.EntityChangeLogMetricsSource;
 import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
 import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
 import org.apache.gravitino.storage.relational.po.cache.OperateType;
@@ -60,12 +61,17 @@ public class TestEntityChangeLogPoller {
     try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
       mockSessionUtils(sessionUtils, mapper);
 
-      EntityChangeLogPoller poller = new EntityChangeLogPoller(1);
+      EntityChangeLogMetricsSource metrics = new 
EntityChangeLogMetricsSource();
+      EntityChangeLogPoller poller = new EntityChangeLogPoller(1, metrics);
       poller.registerListener(firstListenerRecords::addAll);
       poller.registerListener(secondListenerRecords::addAll);
 
       poller.pollChanges();
       poller.pollChanges();
+      Assertions.assertEquals(
+          4, 
metrics.getMetricRegistry().counter("records-delivered-total").getCount());
+      Assertions.assertEquals(
+          4, 
metrics.getMetricRegistry().counter("records-delivered.anonymous-total").getCount());
     }
 
     Assertions.assertEquals(List.of(first, second), firstListenerRecords);
@@ -201,10 +207,186 @@ public class TestEntityChangeLogPoller {
     try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
       mockSessionUtils(sessionUtils, mapper);
 
-      EntityChangeLogPoller poller = new EntityChangeLogPoller(1);
+      EntityChangeLogMetricsSource metrics = new 
EntityChangeLogMetricsSource();
+      EntityChangeLogPoller poller = new EntityChangeLogPoller(1, metrics);
 
       Assertions.assertDoesNotThrow(poller::pollChanges);
+      Assertions.assertEquals(
+          1, 
metrics.getMetricRegistry().counter("poll-failures-total").getCount());
+      Assertions.assertEquals(
+          0, 
metrics.getMetricRegistry().counter("records-fetched-total").getCount());
+    }
+  }
+
+  @Test
+  void testTailSampleFailureDoesNotSuppressFetchedBatch() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    EntityChangeRecord change = change(1L, "CATALOG", "ml1.cat1");
+    when(mapper.selectEntityChanges(0L, MAX_ROWS)).thenReturn(List.of(change));
+    when(mapper.selectMaxChangeId()).thenThrow(new RuntimeException("tail 
query failed"));
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    List<EntityChangeRecord> received = new ArrayList<>();
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      EntityChangeLogPoller poller = new EntityChangeLogPoller(1, metrics);
+      poller.registerListener(received::addAll);
+      poller.pollChanges();
+    }
+
+    Assertions.assertEquals(List.of(change), received);
+    Assertions.assertEquals(
+        1L, 
metrics.getMetricRegistry().getGauges().get("cursor-id").getValue());
+    // The tail was never sampled: the gauge stays at its initial value and 
the cursor moves past
+    // it, which clamps record-lag to zero. The freshness gauge is what 
reveals that state.
+    Assertions.assertEquals(
+        0L, 
metrics.getMetricRegistry().getGauges().get("db-tail-id").getValue());
+    Assertions.assertEquals(
+        0L, 
metrics.getMetricRegistry().getGauges().get("record-lag").getValue());
+    Assertions.assertEquals(
+        -1L,
+        metrics
+            .getMetricRegistry()
+            .getGauges()
+            .get("seconds-since-last-successful-tail-sample")
+            .getValue());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("records-fetched-total").getCount());
+    Assertions.assertEquals(
+        0, 
metrics.getMetricRegistry().counter("poll-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("tail-sample-failures-total").getCount());
+  }
+
+  @Test
+  void testTailSampleFailureRetainsPreviousTail() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    when(mapper.selectEntityChanges(0L, MAX_ROWS))
+        .thenReturn(List.of(change(1L, "CATALOG", "ml1.cat1")));
+    when(mapper.selectMaxChangeId()).thenThrow(new RuntimeException("tail 
query failed"));
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    metrics.setDbTailId(9L);
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      new EntityChangeLogPoller(1, metrics).pollChanges();
+    }
+
+    Assertions.assertEquals(
+        9L, 
metrics.getMetricRegistry().getGauges().get("db-tail-id").getValue());
+    Assertions.assertEquals(
+        1L, 
metrics.getMetricRegistry().getGauges().get("cursor-id").getValue());
+    Assertions.assertEquals(
+        8L, 
metrics.getMetricRegistry().getGauges().get("record-lag").getValue());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("tail-sample-failures-total").getCount());
+  }
+
+  @Test
+  void testInterruptedPollIsNotCountedAsFailure() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    when(mapper.selectEntityChanges(0L, MAX_ROWS))
+        .thenThrow(new RuntimeException(new InterruptedException("shutdown")));
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      new EntityChangeLogPoller(1, metrics).pollChanges();
+      Assertions.assertTrue(Thread.currentThread().isInterrupted());
+      Assertions.assertEquals(
+          0, 
metrics.getMetricRegistry().counter("poll-failures-total").getCount());
+      Assertions.assertEquals(
+          0, 
metrics.getMetricRegistry().counter("tail-sample-failures-total").getCount());
+    } finally {
+      Thread.interrupted();
+    }
+  }
+
+  @Test
+  void testInterruptedTailSampleStopsDeliveryWithoutCountingFailure() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    when(mapper.selectEntityChanges(0L, MAX_ROWS))
+        .thenReturn(List.of(change(1L, "CATALOG", "ml1.cat1")));
+    when(mapper.selectMaxChangeId())
+        .thenThrow(new RuntimeException(new InterruptedException("shutdown")));
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+    List<EntityChangeRecord> received = new ArrayList<>();
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      EntityChangeLogPoller poller = new EntityChangeLogPoller(1, metrics);
+      poller.registerListener(received::addAll);
+      poller.pollChanges();
+      Assertions.assertTrue(Thread.currentThread().isInterrupted());
+      Assertions.assertTrue(received.isEmpty());
+      Assertions.assertEquals(
+          0L, 
metrics.getMetricRegistry().getGauges().get("cursor-id").getValue());
+      Assertions.assertEquals(
+          0, 
metrics.getMetricRegistry().counter("poll-failures-total").getCount());
+      Assertions.assertEquals(
+          0, 
metrics.getMetricRegistry().counter("tail-sample-failures-total").getCount());
+    } finally {
+      Thread.interrupted();
+    }
+  }
+
+  @Test
+  void testSuccessfulEmptyPollSamplesTailAndUpdatesMetrics() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    when(mapper.selectEntityChanges(0L, MAX_ROWS)).thenReturn(List.of());
+    when(mapper.selectMaxChangeId()).thenReturn(4L);
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      new EntityChangeLogPoller(1, metrics).pollChanges();
     }
+
+    Assertions.assertEquals(
+        4L, 
metrics.getMetricRegistry().getGauges().get("db-tail-id").getValue());
+    Assertions.assertEquals(
+        4L, 
metrics.getMetricRegistry().getGauges().get("record-lag").getValue());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().histogram("batch-size-records").getCount());
+    Assertions.assertEquals(
+        0, 
metrics.getMetricRegistry().counter("records-fetched-total").getCount());
+    Assertions.assertTrue(
+        ((Number)
+                    metrics
+                        .getMetricRegistry()
+                        .getGauges()
+                        .get("seconds-since-last-successful-poll")
+                        .getValue())
+                .longValue()
+            >= 0);
+  }
+
+  @Test
+  void testListenerFailureIsAttributedAndCursorStillAdvances() {
+    EntityChangeLogMapper mapper = mock(EntityChangeLogMapper.class);
+    when(mapper.selectEntityChanges(0L, MAX_ROWS))
+        .thenReturn(List.of(change(1L, "TABLE", "ml1.cat1.schema1.table1")));
+    when(mapper.selectMaxChangeId()).thenReturn(1L);
+    EntityChangeLogMetricsSource metrics = new EntityChangeLogMetricsSource();
+
+    try (MockedStatic<SessionUtils> sessionUtils = 
mockStatic(SessionUtils.class)) {
+      mockSessionUtils(sessionUtils, mapper);
+      EntityChangeLogPoller poller = new EntityChangeLogPoller(1, metrics);
+      poller.registerListener(
+          changes -> {
+            throw new IllegalStateException("failure");
+          });
+      poller.pollChanges();
+    }
+
+    Assertions.assertEquals(
+        1L, 
metrics.getMetricRegistry().getGauges().get("cursor-id").getValue());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("listener-failures-total").getCount());
+    Assertions.assertEquals(
+        1, 
metrics.getMetricRegistry().counter("listener-failures.anonymous-total").getCount());
+    Assertions.assertEquals(
+        0, 
metrics.getMetricRegistry().counter("records-applied-total").getCount());
   }
 
   @Test
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
index 77a8fd419f..0ac25beede 100644
--- 
a/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/TestRelationalEntityStoreHierarchicalCache.java
@@ -40,10 +40,13 @@ import org.apache.gravitino.meta.CatalogEntity;
 import org.apache.gravitino.meta.SchemaEntity;
 import org.apache.gravitino.meta.SchemaVersion;
 import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.metrics.MetricsSystem;
+import org.apache.gravitino.metrics.source.MetricsSource;
 import org.apache.gravitino.storage.RandomIdGenerator;
 import org.apache.gravitino.utils.HierarchicalSchemaUtil;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.ValueSource;
 import org.mockito.Mockito;
@@ -63,6 +66,8 @@ public class TestRelationalEntityStoreHierarchicalCache {
   private RelationalEntityStore store;
   private String dbPath;
   private Object previousConfig;
+  private Object previousMetricsSystem;
+  private boolean replacedMetricsSystem;
 
   @AfterEach
   void tearDown() throws Exception {
@@ -75,6 +80,42 @@ public class TestRelationalEntityStoreHierarchicalCache {
       dbPath = null;
     }
     FieldUtils.writeField(GravitinoEnv.getInstance(), "config", 
previousConfig, true);
+    if (replacedMetricsSystem) {
+      FieldUtils.writeField(
+          GravitinoEnv.getInstance(), "metricsSystem", previousMetricsSystem, 
true);
+    }
+  }
+
+  @Test
+  void testChangeLogMetricsRegistrationAndUnregistration() throws Exception {
+    MetricsSystem metricsSystem = new MetricsSystem();
+    replaceMetricsSystem(metricsSystem);
+    initStore(":");
+
+    MetricsSource source = metricsSystem.getMetricsSource("entity-change-log");
+    Assertions.assertNotNull(source);
+    Assertions.assertNotNull(
+        
metricsSystem.getMetricRegistry().getGauges().get("entity-change-log.record-lag"));
+
+    store.close();
+    store = null;
+    Assertions.assertNull(metricsSystem.getMetricsSource("entity-change-log"));
+    Assertions.assertFalse(
+        
metricsSystem.getMetricRegistry().getMetrics().containsKey("entity-change-log.record-lag"));
+  }
+
+  @Test
+  void testChangeLogMetricsWithoutMetricsSystem() throws Exception {
+    replaceMetricsSystem(null);
+    initStore(":");
+
+    Assertions.assertNotNull(FieldUtils.readField(store, "changeLogMetrics", 
true));
+  }
+
+  private void replaceMetricsSystem(MetricsSystem metricsSystem) throws 
IllegalAccessException {
+    previousMetricsSystem = FieldUtils.readField(GravitinoEnv.getInstance(), 
"metricsSystem", true);
+    FieldUtils.writeField(GravitinoEnv.getInstance(), "metricsSystem", 
metricsSystem, true);
+    replacedMetricsSystem = true;
   }
 
   @ParameterizedTest
diff --git a/design-docs/cache-improvement-design.md 
b/design-docs/cache-improvement-design.md
index 750b04e923..92ce2e1524 100644
--- a/design-docs/cache-improvement-design.md
+++ b/design-docs/cache-improvement-design.md
@@ -19,6 +19,12 @@
 
 # Gravitino Cache Improvement Design
 
+This document records an earlier design proposal. Its change-log polling 
pseudocode and
+one-second defaults do not describe the current implementation. For current 
behavior and
+configuration, see [Multi-Node Support for the Entity Store 
Cache](gravitino-entity-cache-multinode-design.md),
+[Change Log 
Propagation](../docs/gravitino-server-config.md#change-log-propagation), and
+[Entity Change Log Metrics](../docs/metrics.md#entity-change-log-metrics).
+
 ---
 
 ## 1. Background
diff --git a/design-docs/gravitino-entity-cache-multinode-design.md 
b/design-docs/gravitino-entity-cache-multinode-design.md
index 949d85530f..a4dc346642 100644
--- a/design-docs/gravitino-entity-cache-multinode-design.md
+++ b/design-docs/gravitino-entity-cache-multinode-design.md
@@ -28,11 +28,15 @@ date: "2026-07-07"
 Gravitino has two caches:
 
 - The **jcasbin authorization cache** already works with more than one node.
-- The **entity store cache** does not. When a change happens on node A, only 
node A clears its cache. Node B keeps serving the old data until its entry 
expires.
+- The **entity store cache** originally had no cross-node invalidation. A 
change on node A only cleared A's local cache, leaving B's entry until expiry.
 
-Because of this, the only safe way to run more than one node today is to turn 
the entity store cache off (`gravitino.cache.enabled=false`). That is bad for 
read-heavy catalogs, especially Iceberg.
+That limitation motivated the change-log listener described here. The local 
Caffeine cache now
+uses it to invalidate entries on peer nodes.
 
-This document proposes a design to make the entity store cache correct when 
running more than one node. It is a design proposal; no behavior has changed 
yet.
+This document records the design and its rationale. Some sections describe the 
baseline at the
+time of the proposal and later implementation phases. For current 
configuration and monitoring,
+see [Change Log 
Propagation](../docs/gravitino-server-config.md#change-log-propagation) and
+[Entity Change Log Metrics](../docs/metrics.md#entity-change-log-metrics).
 
 ## Goals and Non-Goals
 
@@ -49,7 +53,7 @@ This document proposes a design to make the entity store 
cache correct when runn
 
 ---
 
-## Current Cache Implementation
+## Cache Implementation at the Time of the Proposal
 
 `EntityCache` is already an SPI, chosen by `gravitino.cache.impl` and created 
by `CacheFactory`. There is one implementation today, `CaffeineEntityCache` 
(`caffeine`): an in-memory cache, **one copy per node**, with two kinds of 
entries:
 
@@ -259,7 +263,7 @@ The entity store cache becomes a third consumer of the 
poller, next to the catal
 
 ### Consistency
 
-The cache never holds the truth. Every write goes to the DB first, under the 
version lock, so a stale cache can **never** cause a lost update or a bad 
write. The only thing that can go wrong is that a read on **another node** 
returns an old value for a short time — at most one poll interval, until that 
node drops the key.
+The cache never holds the truth. Every write goes to the DB first, under the 
version lock, so a stale cache can **never** cause a lost update or a bad 
write. A read on **another node** can return an old value until that node 
successfully processes the change log. Under normal operation, that happens on 
the next poll; failures can extend the window.
 
 So the real question is simple: when another node reads an old value, does it 
matter? We went through every alter and every drop, for every cached entity, 
one at a time. A stale read falls into one of two buckets:
 
@@ -308,7 +312,7 @@ A stale read is only a problem for a **per-node** cache, 
and only for the load-b
 
 - **Shared cache (redis, `SHARED`): cache everything.** There is no per-node 
window, so model, model version, Semantic Model, and function are cached like 
any other entity, with no extra work. The writing node clears the one shared 
copy, and every node sees it at once.
 - **Per-node cache (caffeine, `LOCAL_PER_NODE`): do not cache model, model 
version, Semantic Model, or function.** Each holds load-bearing content (a 
version URI, the latest version, a Semantic Model definition, or a function 
implementation) that would be silently wrong on another node during the poll 
window. They are read rarely, so reading them from the DB every time costs 
little, and it keeps the rule simple — an entity type is either in or out, with 
no special per-read handling. If t [...]
-- **Metalake on/off flag — cached like the rest of the metalake.** Disabling 
or deleting a metalake is a rare, tenant-level admin action, so we accept the 
small window instead of adding special handling: the metalake is cached and 
invalidated across nodes through the change log like any other entity, so after 
a disable, another node stops allowing operations within one poll interval.
+- **Metalake on/off flag — cached like the rest of the metalake.** Disabling 
or deleting a metalake is a rare, tenant-level admin action, so we accept the 
propagation window instead of adding special handling: the metalake is cached 
and invalidated across nodes through the change log like any other entity. 
After a disable, another node stops allowing operations when it processes the 
change.
 
 We do **not** need a per-entity version check on the cache: for a point read, 
checking the DB version costs the same query as just reading the row, so it 
would buy nothing.
 
@@ -327,17 +331,17 @@ We do **not** need a per-entity version check on the 
cache: for a point read, ch
 | function             | the cached value *is* the code that runs, so a stale 
copy would run the wrong code                | **not cached — read from the 
DB** (revisit if it gets hot)             | cache            |
 | user, group, role    | derived fields (`roleNames`, `securableObjects`) need 
a reverse lookup                            | not cached                        
                                     | not cached       |
 
-After this, everything a per-node cache serves is safe or bounded: 
connector-backed entities are safe by construction; self-contained store 
entities are only ever cosmetically stale (and a rare metalake disable 
self-corrects within one poll interval); model / model version / Semantic Model 
/ function are read from the DB. A shared cache is safe throughout because it 
has no window.
+After this, everything a per-node cache serves is safe or bounded by 
change-log processing under normal operation: connector-backed entities are 
safe by construction; self-contained store entities are only ever cosmetically 
stale (and a rare metalake disable self-corrects after propagation); model / 
model version / Semantic Model / function are read from the DB. A shared cache 
is safe throughout because it has no window.
 
-#### The staleness promise (SLA)
+#### Expected propagation and failure behavior
 
 For everything that stays in the cache:
 
 - The node that made the change sees it right away.
-- Every other node sees it within **one poll interval**. This is set by 
`gravitino.entityChangeLog.pollIntervalSecs` (**default 3 seconds**; lower it, 
e.g. to 1 second, for a shorter delay at the cost of more frequent DB polls).
+- Under normal operation, every other node sees it after the next successful 
poll. The interval is set by `gravitino.entityChangeLog.pollIntervalSecs` 
(**default 3 seconds**; lower it, e.g. to 1 second, for a shorter delay at the 
cost of more frequent DB polls). Database failures or a listener that fails to 
recover can extend this delay.
 - The cache's own TTL (minutes to hours) is only a safety net in case the 
poller ever misses a row; it is not the main mechanism.
 
-This "at most one poll interval" promise holds only if the poller never drops 
a row. So the poller must reuse the same gap-safe, id-based polling already 
built for the change log (see `#11736`), keep the TTL as a backstop, and be 
watched for lag.
+The poller uses gap-safe, id-based polling (see `#11736`). It delivers a batch 
once and advances its shared cursor even when a listener fails, so each 
listener must recover locally by clearing its cache. The cache TTL remains a 
backstop. Operators can use the [entity change log 
metrics](../docs/metrics.md#entity-change-log-metrics) to watch the sampled 
database tail, cursor, record lag, last successful poll age, listener failures, 
and fallback clears; the debug logs correlate the write,  [...]
 
 ---
 
@@ -398,13 +402,13 @@ gravitino.cache.redis.serializer = ...         # e.g. 
JSON or a binary codec
 
 ## Choosing an Implementation
 
-|              | `caffeine` (default)                | `redis` (optional)      
                        |
-| ------------ | ----------------------------------- | 
----------------------------------------------- |
-| Dependency   | none                                | a Redis the operator 
runs                       |
-| Read latency | local memory                        | one network round-trip  
                        |
-| Consistency  | eventual (≤ one poll interval)      | strong 
(read-your-writes)                       |
-| Transport    | reuses `entity_change_log` + poller | none — one shared copy  
                        |
-| Best for     | most deployments                    | already running Redis; 
wants strong consistency |
+|              | `caffeine` (default)                                   | 
`redis` (optional)                              |
+| ------------ | ------------------------------------------------------ | 
----------------------------------------------- |
+| Dependency   | none                                                   | a 
Redis the operator runs                       |
+| Read latency | local memory                                           | one 
network round-trip                          |
+| Consistency  | eventual (next successful poll under normal operation) | 
strong (read-your-writes)                       |
+| Transport    | reuses `entity_change_log` + poller                    | none 
— one shared copy                          |
+| Best for     | most deployments                                       | 
already running Redis; wants strong consistency |
 
 Both are chosen through the same SPI, so a user picks by environment with a 
single config change. `caffeine` stays the default.
 
@@ -421,16 +425,16 @@ Both are chosen through the same SPI, so a user picks by 
environment with a sing
 
 ## Test Plan
 
-| Area                    | Check                                              
                                                                                
                                                                                
          |
-| ----------------------- | 
----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 |
-| Multi-node — caffeine   | node A runs ALTER/DROP on table / schema / 
catalog; node B serves the fresh entity within one poll interval                
                                                                                
                  |
-| Mutation-event coverage | create emits nothing; 
overwrite/update/rename/successful drop emit exactly one row for every 
cacheable type; rename keeps only the old key; entity and row roll back 
together                                                 |
-| Tag / policy cross-node | an ALTER/DROP of a tag or policy on node A is 
reflected on node B within one poll interval                                    
                                                                                
                |
-| Hierarchical drop       | dropping a schema drops the schema's cached child 
tables on the other node(s), and leaves a sibling schema alone (forward prefix 
scan)                                                                           
            |
-| Shared feed             | the entity store cache, the catalog cache, and the 
jcasbin id-mapping cache all act on the same structural rows; adding the entity 
store consumer does not change the others                                       
          |
-| Not cached (caffeine)   | get user / group / role / model / model version / 
function return correct data straight from the DB; authorization is unaffected  
                                                                                
           |
-| Silent-staleness guard  | disabling a metalake on A blocks operations on B 
within one poll interval (the metalake is cached and invalidated cross-node); 
under redis, model / model version / function are cached and always current 
(one shared copy) |
-| Multi-node — redis      | node A ALTER/DROP; node B reads the fresh entity 
right away; a container drop removes child keys via `ZRANGEBYLEX`; no half-done 
drop is visible                                                                 
            |
-| Redis stale-write       | a stale `v1` write after a committed `v2` + delete 
is rejected by the version guard; no node serves a value older than the last 
commit                                                                          
             |
-| Relation reads          | owner / role / tag / policy listings return 
correct results from the DB with no caching; tag/policy inheritance still 
resolves                                                                        
                       |
-| Regression              | single-node behavior, the write path's version 
check, and `list` strong consistency are unchanged                              
                                                                                
              |
+| Area                    | Check                                              
                                                                                
                                              |
+| ----------------------- | 
--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
 |
+| Multi-node — caffeine   | node A runs ALTER/DROP on table / schema / 
catalog; node B serves the fresh entity after the next successful poll          
                                                      |
+| Mutation-event coverage | create emits nothing; 
overwrite/update/rename/successful drop emit exactly one row for every 
cacheable type; rename keeps only the old key; entity and row roll back 
together    |
+| Tag / policy cross-node | an ALTER/DROP of a tag or policy on node A is 
reflected on node B after the next successful poll                              
                                                   |
+| Hierarchical drop       | dropping a schema drops the schema's cached child 
tables on the other node(s), and leaves a sibling schema alone (forward prefix 
scan)                                           |
+| Shared feed             | the entity store cache, the catalog cache, and the 
jcasbin id-mapping cache all act on the same structural rows; adding the entity 
store consumer does not change the others     |
+| Not cached (caffeine)   | get user / group / role / model / model version / 
function return correct data straight from the DB; authorization is unaffected  
                                               |
+| Silent-staleness guard  | disabling a metalake on A blocks operations on B 
after B processes the change log; under redis, model / model version / function 
are cached and always current (one shared copy) |
+| Multi-node — redis      | node A ALTER/DROP; node B reads the fresh entity 
right away; a container drop removes child keys via `ZRANGEBYLEX`; no half-done 
drop is visible                                 |
+| Redis stale-write       | a stale `v1` write after a committed `v2` + delete 
is rejected by the version guard; no node serves a value older than the last 
commit                                           |
+| Relation reads          | owner / role / tag / policy listings return 
correct results from the DB with no caching; tag/policy inheritance still 
resolves                                                   |
+| Regression              | single-node behavior, the write path's version 
check, and `list` strong consistency are unchanged                              
                                                  |
diff --git a/docs/gravitino-server-config.md b/docs/gravitino-server-config.md
index 1fb42dea4a..965e1a71c8 100644
--- a/docs/gravitino-server-config.md
+++ b/docs/gravitino-server-config.md
@@ -141,10 +141,13 @@ it recognizes.
 ### Running More Than One Server
 
 Servers behind a load balancer share the entity store but keep local caches. 
Each server polls the
-entity change log and invalidates entries that another server has modified. 
The defaults are safe:
-a three second poll, and a server that cannot keep its caches current exits 
rather than serving
-metadata it knows to be stale. Point the load balancer's health check at `GET 
/health/ready` so a
-server that has lost its database stops receiving traffic.
+entity change log and invalidates entries that another server has modified. 
The default poll
+interval is three seconds. The poller delivers each batch to every registered 
listener once and
+then advances its cursor; a listener that cannot invalidate a key must clear 
its local cache. The
+poller logs query and listener failures and continues polling. Monitor the
+[entity change log metrics](metrics.md#entity-change-log-metrics), especially 
record lag, time
+since the last successful poll, listener failures, and fallback clears. Point 
the load balancer's
+health check at `GET /health/ready` so a server that has lost its database 
stops receiving traffic.
 
 Jobs run by the default `local` job executor keep their output in 
`gravitino.job.stagingDir`. Put
 that directory on storage shared by all servers, for example an NFS mount, so 
that a request for a
diff --git a/docs/metrics.md b/docs/metrics.md
index 9102191999..b7e19b32c9 100644
--- a/docs/metrics.md
+++ b/docs/metrics.md
@@ -71,3 +71,43 @@ 
gravitino_catalog_datasource_idle_connections{provider="jdbc",metalake="test_met
 
gravitino_catalog_datasource_active_connections{provider="jdbc",metalake="test_metalake",catalog="test_catalog",}
 0.0
 
gravitino_catalog_datasource_max_connections{provider="jdbc",metalake="test_metalake",catalog="test_catalog",}
 10.0
 ```
+
+#### Entity Change Log Metrics
+
+The `entity-change-log` source exposes each server's change-log processing 
state through JMX and
+`/prometheus/metrics`. For example, `entity-change-log.record-lag` in the 
metrics registry becomes
+`entity_change_log_record_lag` in Prometheus. Gauges read only in-memory 
values; the poller samples
+the database tail once per cycle. If only the tail sample fails, delivery 
continues and the tail
+value remains at its last successful sample.
+
+| Metric suffix                                                | Type and unit 
           | Meaning                                                            
                                                                                
         |
+| ------------------------------------------------------------ | 
------------------------ | 
-----------------------------------------------------------------------------------------------------------------------------------------------------------
 |
+| `db-tail-id`, `cursor-id`                                    | gauge, change 
ID         | Latest sampled database ID and last delivered ID on this server.   
                                                                                
         |
+| `record-lag`                                                 | gauge, 
records           | Sampled tail minus cursor, clamped at zero. Interpret only 
while tail sampling succeeds.                                                   
                 |
+| `seconds-since-last-successful-poll`                         | gauge, 
seconds           | Time since a successful database poll, including an empty 
result; `-1` before the first poll.                                             
                  |
+| `seconds-since-last-successful-tail-sample`                  | gauge, 
seconds           | Time since `db-tail-id` was last refreshed; `-1` before the 
first sample. Trust `db-tail-id` and `record-lag` only while this stays near 
the poll interval. |
+| `poll-failures-total`                                        | counter, 
failures        | Failed poll cycles.                                           
                                                                                
              |
+| `tail-sample-failures-total`                                 | counter, 
failures        | Failed database-tail samples; fetched batches can still be 
delivered.                                                                      
                 |
+| `listener-failures-total`, `listener-failures.<class>-total` | counter, 
failures        | Total failures and failures by registered listener class.     
                                                                                
              |
+| `records-fetched-total`, `records-delivered-total`           | counter, 
records         | Rows fetched and rows delivered successfully to listeners; 
one row delivered to two listeners counts twice as delivered.                   
                 |
+| `records-delivered.<class>-total`                            | counter, 
records         | Successful deliveries by listener class. Lambda and anonymous 
listeners share the `anonymous` bucket.                                         
              |
+| `records-applied-total`                                      | counter, 
invalidations   | Targeted entity-cache invalidations completed successfully; 
malformed rows and fallback clears do not count.                                
                |
+| `batch-size-records`                                         | histogram, 
records       | Number of rows fetched per successful poll, including empty 
polls.                                                                          
                |
+| `poll-duration`                                              | timer, 
duration          | End-to-end poll-cycle duration.                             
                                                                                
                |
+| `invalidation-failures-total`, `fallback-clears-total`       | counter, 
failures/clears | Failed targeted entity-cache invalidations and successful 
full-cache recovery clears.                                                     
                  |
+
+The poller delivers each batch once and has no pending or retry state. A 
failed listener must
+recover locally; its failure counter and log identify the affected listener. 
The debug logs use
+`entityChangeLog` fields to trace an append, poll, delivery, and invalidation. 
Append logs mean the
+row was added to the current transaction, not that the transaction committed.
+
+For an incident, check `seconds-since-last-successful-tail-sample` before 
comparing `db-tail-id`
+with `cursor-id` on the affected server. If it exceeds the poll interval, the 
tail sample itself is
+failing: the retained tail can fall below an advancing cursor and `record-lag` 
can read zero despite
+an unknown database tail. `tail-sample-failures-total` counts those failures 
for alerting. With a
+fresh tail sample, a growing `record-lag` together with an increasing poll age 
or
+`poll-failures-total` points to polling trouble.
+If the cursor advances but data remains stale, inspect 
`listener-failures-total`, `records-delivered.<class>-total`,
+`invalidation-failures-total`, and `fallback-clears-total`, then correlate the 
debug logs by
+encoded `fullName` and change ID. The sampled tail and cursor are 
process-local; each server has
+its own values.

Reply via email to