yuqi1129 commented on code in PR #13388:
URL: https://github.com/apache/gravitino/pull/13388#discussion_r4068896098
##########
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:
Fixed in 182b832048. The tail query is now best effort: a normal sampling
failure keeps the previous gauge value, records the fetch as successful, and
delivers the fetched batch. An interruption still stops the cycle. Added
regression tests for both paths.
##########
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:
Fixed in 182b832048. The interrupt check now runs before pollFailed(), with
tests for an interrupted fetch and an interrupted tail sample.
##########
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:
Fixed in 182b832048. Named listeners keep their class name; synthetic lambda
and anonymous listener classes use one stable `anonymous` metric bucket. The
poller test verifies that two lambda listeners aggregate there.
##########
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:
I kept the single-argument constructors because both existed before this PR
and removing them would break existing callers. The cache listener now always
has a non-null metrics source, so its hot path no longer needs nullable guards.
The legacy constructors are documented as using an unregistered source;
production passes the shared registered source.
##########
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:
Added `records-delivered.<class>-total` alongside the aggregate counter in
182b832048, using the same stable listener name as the failure counter. Updated
the metrics docs and tests, including the aggregate `anonymous` bucket for
lambda listeners.
--
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]