jerryshao commented on code in PR #13388:
URL: https://github.com/apache/gravitino/pull/13388#discussion_r4068328273
##########
core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java:
##########
@@ -198,8 +215,21 @@ private synchronized void doPollChanges() {
@Nullable
private BatchDelivery fetchNextDelivery() {
+ long fetchStartNanos = System.nanoTime();
List<EntityChangeRecord> changes = fetchEntityChanges();
+ long dbTailId =
+ getOrDefault(
+ SessionUtils.getWithoutCommit(
+ EntityChangeLogMapper.class,
EntityChangeLogMapper::selectMaxChangeId));
+ metrics.setDbTailId(dbTailId);
+ metrics.pollSucceeded(changes.size());
+ long durationMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() -
fetchStartNanos);
Review Comment:
[Important] This observability-only query now sits between the fetch and the
delivery, so a failure here throws away a batch that was already fetched
successfully.
`doPollChanges()` (line 209) calls `fetchNextDelivery()` first and reaches
`deliver()` only if it returns (lines 210-213), so an exception from
`selectMaxChangeId` propagates past `deliver()` into the catch at line 195. The
listeners get nothing that cycle, `poll-failures-total` increments, and the
cursor stays where it was. A transient error costs one poll interval, since
`advanceCursor` (line 312) runs only after delivery and the batch is refetched.
But a query that fails on its own — a statement timeout, or a connection
dropped between the two `SessionUtils.getWithoutCommit` calls, each of which
opens and closes its own session — stalls cross-node cache invalidation for as
long as it keeps failing, while the fetch that actually matters is working
fine. It also doubles the change-log query count per poll per node (every 3s by
default).
Keep the tail sample off the critical path: either wrap lines 220-224 in
their own try/catch that only logs and leaves the gauge stale, or skip the
query when `changes.size() < ENTITY_CHANGE_POLLER_MAX_ROWS`, where the tail at
fetch time is already known from the last fetched id (or from the cursor, for
an empty batch).
Verified by: read EntityChangeLogPoller.java:191-247 and 307-315 on this
checkout — `deliver()` is unreachable when `fetchNextDelivery` throws, and
`metrics.setCursorId` is only called from `advanceCursor`, so the cursor cannot
advance past an undelivered batch.
##########
core/src/main/java/org/apache/gravitino/storage/relational/EntityCacheChangeLogListener.java:
##########
@@ -62,70 +65,128 @@ public class EntityCacheChangeLogListener implements
EntityChangeLogListener {
private static final Logger LOG =
LoggerFactory.getLogger(EntityCacheChangeLogListener.class);
private final EntityCache cache;
+ @Nullable private final EntityChangeLogMetricsSource metrics;
/**
* Creates a listener that invalidates the given entity store cache.
*
* @param cache the per-node entity store cache to keep coherent
*/
public EntityCacheChangeLogListener(EntityCache cache) {
+ this(cache, null);
+ }
+
+ /**
+ * Creates a listener with metrics shared with the poller.
+ *
+ * @param cache the per-node entity store cache
+ * @param metrics process-local change log metrics, or null when metrics are
unavailable
+ */
+ public EntityCacheChangeLogListener(
+ EntityCache cache, @Nullable EntityChangeLogMetricsSource metrics) {
Review Comment:
[Nit] This nullable-metrics overload exists only for tests, and it buys four
`if (metrics != null)` guards in the hot loop (lines 117-119, 121-123, and the
two on the fallback path). `grep -rn "new EntityCacheChangeLogListener("` shows
the single production caller (RelationalEntityStore.java:140) always passes a
real source; only tests use the one-argument form. The same applies to
`EntityChangeLogPoller(long)` at EntityChangeLogPoller.java:90-92, whose fresh
source is never registered with the `MetricsSystem`, so any metric it records
is invisible.
Making `metrics` required and having the tests pass `new
EntityChangeLogMetricsSource()` (as the updated tests already do) would drop
the `@Nullable` field, all four guards, and a public constructor that silently
discards metrics.
Verified by: grepped every `new EntityCacheChangeLogListener(` and `new
EntityChangeLogPoller(` call site in the repo on this checkout; the only
non-test uses are RelationalEntityStore.java:114 and :140, both passing the
registered source.
##########
core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java:
##########
@@ -174,9 +190,10 @@ public void close() {
@VisibleForTesting
void pollChanges() {
- try {
+ try (Timer.Context ignored = metrics.timePoll()) {
doPollChanges();
} catch (Throwable e) {
+ metrics.pollFailed();
Review Comment:
[Nit] `pollFailed()` runs before the interrupt check at line 202, so a poll
interrupted during shutdown is counted as a real poll failure. `close()` calls
`scheduler.shutdownNow()` once the 5-second `awaitTermination` expires (lines
177-179, and again at 181), which interrupts an in-flight poll;
`handleInterruptIfAny` then recognises it and returns without even logging a
failure — but the counter has already moved. docs/metrics.md points operators
at `poll-failures-total` as a polling-trouble signal, so moving
`metrics.pollFailed()` below the interrupt check keeps that counter to failures
that really are failures.
Verified by: read EntityChangeLogPoller.java:191-207 and close() at 173-189
on this checkout.
##########
core/src/main/java/org/apache/gravitino/metrics/source/EntityChangeLogMetricsSource.java:
##########
@@ -0,0 +1,109 @@
+/*
+ * 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 Counter pollFailures = getCounter("poll-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>)
+ () -> {
+ long last = lastSuccessfulPollMs.get();
+ return last == 0 ? -1 : Math.max(0, (System.currentTimeMillis()
- last) / 1000);
+ });
+ }
+
+ /** Records the database tail sampled by a poll, without querying from the
gauge. */
+ public void setDbTailId(long id) {
+ dbTailId.set(id);
+ }
+
+ /** 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 listener delivery that failed, attributed by its stable class
name. */
+ public void listenerFailed(String listenerName) {
+ listenerFailures.inc();
+ getCounter("listener-failures." + listenerName.replace('.', '_') +
"-total").inc();
Review Comment:
[Nit] The Javadoc above says "stable class name", but the name is only
stable for named classes. `EntityChangeLogListener` is annotated
`@FunctionalInterface` (EntityChangeLogListener.java:25) and
`registerEntityChangeLogListener` is public API on the store, so a lambda or
anonymous listener is an expected shape — the poller's own tests register
lambdas (TestEntityChangeLogPoller.java:94, 133, 259). For a lambda,
`getClass().getName()` looks like
`...TestEntityChangeLogPoller$$Lambda$14/0x000000080019c440`, which changes
between JVM runs, so each restart creates a new metric name in a registry that
never drops entries, and the resulting Prometheus series cannot be charted
across restarts.
The three listeners shipped today are named classes, so nothing is broken
right now. Guarding it cheaply would help: fall back to something stable when
the name contains `$$Lambda` (for example `listener.getClass().getSuperclass()`
or a fixed `anonymous` bucket), or have the listener interface expose a name.
Verified by: read EntityChangeLogMetricsSource.java:79-83 and
EntityChangeLogListener.java:24-26 on this checkout; grepped the three
implementations (`EntityCacheChangeLogListener`, `CatalogChangeLogListener`,
`JcasbinChangeListener`) — all named classes.
##########
core/src/main/java/org/apache/gravitino/storage/relational/EntityChangeLogPoller.java:
##########
@@ -308,13 +341,21 @@ private void notifyListeners(BatchDelivery delivery) {
}
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(delivery.changes.size());
Review Comment:
[Question] `records-delivered-total` is incremented once per listener, so
with the three listeners registered today (`EntityCacheChangeLogListener`,
`CatalogChangeLogListener`, `JcasbinChangeListener`) it runs at roughly 3x
`records-fetched-total`. docs/metrics.md is explicit about this, so it is
clearly deliberate — but was a per-listener breakdown considered here too?
Failures are already attributed per listener
(`listener-failures.<class>-total`), while deliveries are only aggregate, so an
operator cannot tell which listener stopped receiving rows: the natural
`fetched` vs `delivered` comparison is meaningless once the multiplier depends
on how many listeners happen to be registered. A
`records-delivered.<class>-total` alongside the existing per-listener failure
counter would make the two halves symmetric.
Verified by: read EntityChangeLogPoller.java:336-358 (the increment is
inside the per-listener loop) and grepped `implements EntityChangeLogListener`
— three implementations in the repo.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]