jerryshao commented on code in PR #13388: URL: https://github.com/apache/gravitino/pull/13388#discussion_r4069552874
########## core/src/main/java/org/apache/gravitino/metrics/source/EntityChangeLogMetricsSource.java: ########## @@ -0,0 +1,116 @@ +/* + * 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 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>) + () -> { + 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 failed database-tail sample while allowing an already fetched batch to proceed. */ + public void tailSampleFailed() { Review Comment: [Question] Was a freshness signal considered alongside this counter, or is the counter deliberately the whole answer? `tail-sample-failures-total` is monotonic, so it answers "has the tail sample ever failed", not "is the tail I am reading right now stale" - and docs/metrics.md:103 asks operators to check it *before* trusting `db-tail-id` and `record-lag`. Over Prometheus a `rate()` closes that gap. Over JMX, which docs/metrics.md:77 lists as an equal export path, there is no rate: one transient failure hours ago leaves the counter permanently non-zero, and the playbook then reads as "never trust these gauges again". A `seconds-since-last-successful-tail-sample` gauge, shaped exactly like `seconds-since-last-successful-poll` (lines 49-56), would answer the question the playbook actually asks and costs one `AtomicLong`. Not blocking - the counter is a clear improvement on the previous silence. Verified by: read EntityChangeLogMetricsSource.java:29-56 and 80-83, EntityChangeLogPoller.java:226-236, and docs/metrics.md:75-81 and 102-107 on this checkout; the source registers four gauges and none of them tracks tail-sample freshness. ########## core/src/test/java/org/apache/gravitino/storage/relational/TestEntityChangeLogPoller.java: ########## @@ -201,12 +207,156 @@ void testPollChangesCatchesFetchFailures() { 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(); + metrics.setDbTailId(0L); Review Comment: [Nit] Reseeding the tail to 0 makes half of this test tautological, and it leaves the behaviour `182b832` introduced asserted nowhere. `dbTailId` is an `AtomicLong` that already starts at 0 (EntityChangeLogMetricsSource.java:29), so `setDbTailId(0L)` is a no-op: the `db-tail-id == 0` assertion on line 242 passes whether the gauge retained its previous value or was never written at all. The `record-lag == 0` / tail-below-cursor state this reparameterisation buys is the more valuable case and worth keeping. But the thing the best-effort catch was added for - what its own warn message calls "retaining the previous gauge value" - is no longer covered: `grep -rn setDbTailId core/src/test` finds only this line and TestEntityChangeLogMetricsSource.java:33, and that test never fails a sample. Both need two tests, since a cursor that advances past the tail requires the tail to be 0 here: keep this one as is, and add one that seeds a non-zero tail (say 9), fails `selectMaxChangeId`, and asserts `db-tail-id` is still 9 afterwards. Verified by: read TestEntityChangeLogPoller.java:221-251 and EntityChangeLogMetricsSource.java:29 and 44-48 on this checkout; grepped every `setDbTailId` call under `core/src/test`. ########## docs/gravitino-server-config.md: ########## @@ -157,22 +160,22 @@ empty string or list; `(none)` means it has no default at all. #### HTTP Server -| Configuration Item | Description | Default Value | -|------------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|-----------------------------------------| -| `gravitino.server.webserver.host` | The address the server binds to. | `0.0.0.0` | -| `gravitino.server.webserver.httpPort` | The port the server listens on. | `8090` | -| `gravitino.server.webserver.minThreads` | Minimum threads in the Jetty thread pool. Values below 8 are raised to 8. | Twice the processor count, 8 to 100 | -| `gravitino.server.webserver.maxThreads` | Maximum threads in the Jetty thread pool. Values below 8 are raised to 8, and the value must be at least `minThreads`. | Four times the processor count, min 400 | -| `gravitino.server.webserver.threadPoolWorkQueueSize` | Size of the Jetty thread pool work queue. | `100` | -| `gravitino.server.webserver.idleTimeout` | Timeout in milliseconds for idle connections. | `30000` | -| `gravitino.server.webserver.stopTimeout` | Time in milliseconds Jetty waits for a graceful shutdown. See `org.eclipse.jetty.server.Server#setStopTimeout`. | `30000` | -| `gravitino.server.shutdown.timeout` | Time in milliseconds for the Gravitino server itself to shut down gracefully. | `3000` | -| `gravitino.server.webserver.requestHeaderSize` | Maximum size in bytes of an HTTP request header. | `131072` | -| `gravitino.server.webserver.responseHeaderSize` | Maximum size in bytes of an HTTP response header. | `131072` | -| `gravitino.server.webserver.customFilters` | Comma-separated list of servlet filter class names to apply to the API. | (empty) | -| `gravitino.server.rest.extensionPackages` | Comma-separated list of packages to scan for additional REST resources. | (empty) | -| `gravitino.server.visibleConfigs` | Comma-separated list of extra properties to expose on the unauthenticated `GET /configs` endpoint, on top of the fixed set it always returns. Additive, so each entry widens what is public. | (empty) | -| `gravitino.server.bulk.maxItems` | Maximum number of items allowed in a single bulk request. | `100` | +| Configuration Item | Description | Default Value | Review Comment: [Nit] This reformat is outside the PR's subject and could be dropped to keep the diff reviewable. The HTTP Server table (lines 163-178), the KMS table (570-574) and two rows of the environment-variable table (672, 684) have nothing to do with change-log observability; the same commit also repads three tables in the two design docs. That is ~46 changed lines in this file alone that a reviewer has to diff character by character to confirm are inert, and it is the kind of churn that conflicts with any other PR touching these sections. No Markdown linter runs over `docs/` in CI (`.github/workflows` has a prettier step for the web UI only), so this is not a CI requirement either. If the padding is wanted, a separate docs-only PR is the cheaper place for it. As a side note, it also does not fully achieve what the commit message claims: the separator row at line 164 is two characters wider than the header row above it. Verified by: read docs/gravitino-server-config.md:160-178 and 567-574 on this checkout; confirmed the commit is inert by stripping spaces, pipes and dashes from this file at `182b832` and at `HEAD` and diffing - identical; grepped `.github/workflows/*.yml` for `markdownlint`/`prettier`. -- 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]
