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.