This is an automated email from the ASF dual-hosted git repository. voonhous pushed a commit to branch release-1.2.1 in repository https://gitbox.apache.org/repos/asf/hudi.git
commit f1259ca8cd0274aced9dc9d702609f377a3923b1 Author: Peter Huang <[email protected]> AuthorDate: Tue Jun 9 07:42:26 2026 -0700 fix(metrics): NPE handling when hudi metrics is disabled (#18947) (cherry picked from commit 1fd2c3671a8081227727f5d6e175ef49d9a66f2d) --- .../org/apache/hudi/metrics/HoodieMetrics.java | 11 +- .../org/apache/hudi/metrics/TestHoodieMetrics.java | 316 +++++++++++++++++++++ .../sync/common/metrics/HoodieMetaSyncMetrics.java | 3 + 3 files changed, 326 insertions(+), 4 deletions(-) diff --git a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metrics/HoodieMetrics.java b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metrics/HoodieMetrics.java index c777e842bbd3..941085f451bd 100644 --- a/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metrics/HoodieMetrics.java +++ b/hudi-client/hudi-client-common/src/main/java/org/apache/hudi/metrics/HoodieMetrics.java @@ -266,7 +266,7 @@ public class HoodieMetrics { } public Timer.Context getConflictResolutionCtx() { - if (config.isLockingMetricsEnabled() && conflictResolutionTimer == null) { + if (config.isMetricsOn() && config.isLockingMetricsEnabled() && conflictResolutionTimer == null) { conflictResolutionTimer = createTimer(conflictResolutionTimerName); } return conflictResolutionTimer == null ? null : conflictResolutionTimer.time(); @@ -555,6 +555,9 @@ public class HoodieMetrics { * Given a commit action, metrics name and value this method reports custom metrics. */ public void reportMetrics(String commitAction, String metricName, long value) { + if (!config.isMetricsOn()) { + return; + } metrics.registerGauge(getMetricsName(commitAction, metricName), value); } @@ -566,7 +569,7 @@ public class HoodieMetrics { } public void emitConflictResolutionSuccessful() { - if (config.isLockingMetricsEnabled()) { + if (config.isMetricsOn() && config.isLockingMetricsEnabled()) { log.info("Sending conflict resolution success metric"); conflictResolutionSuccessCounter = getCounter(conflictResolutionSuccessCounter, conflictResolutionSuccessCounterName); conflictResolutionSuccessCounter.inc(); @@ -574,7 +577,7 @@ public class HoodieMetrics { } public void emitConflictResolutionFailed() { - if (config.isLockingMetricsEnabled()) { + if (config.isMetricsOn() && config.isLockingMetricsEnabled()) { log.info("Sending conflict resolution failure metric"); conflictResolutionFailureCounter = getCounter(conflictResolutionFailureCounter, conflictResolutionFailureCounterName); conflictResolutionFailureCounter.inc(); @@ -582,7 +585,7 @@ public class HoodieMetrics { } public void emitConflictResolutionByCategory(HoodieWriteConflictException.ConflictCategory category) { - if (config.isLockingMetricsEnabled()) { + if (config.isMetricsOn() && config.isLockingMetricsEnabled()) { switch (category) { case INGESTION_VS_INGESTION: conflictResolutionIngestionVsIngestionCounter = getCounter( diff --git a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metrics/TestHoodieMetrics.java b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metrics/TestHoodieMetrics.java index b785d9025288..1a2702a77a04 100755 --- a/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metrics/TestHoodieMetrics.java +++ b/hudi-client/hudi-client-common/src/test/java/org/apache/hudi/metrics/TestHoodieMetrics.java @@ -40,6 +40,8 @@ import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; import java.util.Arrays; +import java.util.HashMap; +import java.util.Map; import java.util.Random; import java.util.UUID; import java.util.stream.Stream; @@ -48,7 +50,10 @@ import static org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR import static org.apache.hudi.metrics.HoodieMetrics.COUNTER_METRIC_EXTENSION; import static org.apache.hudi.metrics.HoodieMetrics.FAILURE_COUNTER; import static org.apache.hudi.metrics.HoodieMetrics.SOURCE_READ_AND_INDEX_ACTION; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -322,6 +327,317 @@ public class TestHoodieMetrics { } } + // ----------------------------------------------------------------------- + // Metrics-off safety tests (NPE guards) + // ----------------------------------------------------------------------- + + private HoodieMetrics buildMetricsOff() { + HoodieWriteConfig offConfig = mock(HoodieWriteConfig.class); + when(offConfig.isMetricsOn()).thenReturn(false); + when(offConfig.getTableName()).thenReturn("test_table"); + return new HoodieMetrics(offConfig, HoodieTestUtils.getDefaultStorage()); + } + + @Test + public void testTimerContextsReturnNullWhenMetricsOff() { + HoodieMetrics metricsOff = buildMetricsOff(); + assertNull(metricsOff.getRollbackCtx()); + assertNull(metricsOff.getCompactionCtx()); + assertNull(metricsOff.getLogCompactionCtx()); + assertNull(metricsOff.getClusteringCtx()); + assertNull(metricsOff.getCleanCtx()); + assertNull(metricsOff.getArchiveCtx()); + assertNull(metricsOff.getCommitCtx()); + assertNull(metricsOff.getFinalizeCtx()); + assertNull(metricsOff.getDeltaCommitCtx()); + assertNull(metricsOff.getIndexCtx()); + assertNull(metricsOff.getSourceReadAndIndexTimerCtx()); + assertNull(metricsOff.getConflictResolutionCtx()); + } + + @Test + public void testUpdateMethodsAreNoOpsWhenMetricsOff() { + HoodieMetrics metricsOff = buildMetricsOff(); + HoodieCommitMetadata metadata = mock(HoodieCommitMetadata.class); + + assertDoesNotThrow(() -> metricsOff.updateCommitMetrics(0L, 0L, metadata, "commit")); + assertDoesNotThrow(() -> metricsOff.updateRollbackMetrics(0L, 0L)); + assertDoesNotThrow(() -> metricsOff.updateCleanMetrics(0L, 0)); + assertDoesNotThrow(() -> metricsOff.updateFinalizeWriteMetrics(0L, 0L)); + assertDoesNotThrow(() -> metricsOff.updateIndexMetrics("action", 0L)); + assertDoesNotThrow(() -> metricsOff.updateSourceReadAndIndexMetrics("action", 0L)); + assertDoesNotThrow(() -> metricsOff.updateArchiveMetrics(0L, 0)); + assertDoesNotThrow(() -> metricsOff.updateArchivalMetrics(new HashMap<>())); + assertDoesNotThrow(() -> metricsOff.updatePostCommitMetrics(true, 0L)); + assertDoesNotThrow(() -> metricsOff.updatePostCommitMetrics(false, 0L)); + assertDoesNotThrow(() -> metricsOff.reportMetrics("action", "metric", 0L)); + assertDoesNotThrow(() -> metricsOff.updateClusteringFileCreationMetrics(0L)); + assertDoesNotThrow(() -> metricsOff.emitRollbackFailure("SomeException")); + assertDoesNotThrow(() -> metricsOff.emitRollbackFailure(null)); + assertDoesNotThrow(() -> metricsOff.emitCompactionRequested()); + assertDoesNotThrow(() -> metricsOff.emitCompactionCompleted()); + assertDoesNotThrow(() -> metricsOff.emitIndexTypeMetrics(0)); + assertDoesNotThrow(() -> metricsOff.emitMetadataEnablementMetrics(false, false, false, false)); + assertDoesNotThrow(() -> metricsOff.emitVersionMetrics()); + } + + @Test + public void testConflictResolutionMetricsAreNoOpsWhenMetricsOffButLockingEnabled() { + // Key NPE scenario: global metrics are off so metrics field is null. + // The isMetricsOn() && isLockingMetricsEnabled() guard short-circuits before + // reaching getCounter() -> metrics.getRegistry(), so no NPE should occur. + HoodieWriteConfig offConfig = mock(HoodieWriteConfig.class); + when(offConfig.isMetricsOn()).thenReturn(false); + when(offConfig.getTableName()).thenReturn("test_table"); + HoodieMetrics metricsOff = new HoodieMetrics(offConfig, HoodieTestUtils.getDefaultStorage()); + + assertNull(metricsOff.getConflictResolutionCtx()); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionSuccessful()); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionFailed()); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionByCategory( + HoodieWriteConflictException.ConflictCategory.INGESTION_VS_INGESTION)); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionByCategory( + HoodieWriteConflictException.ConflictCategory.INGESTION_VS_TABLE_SERVICE)); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionByCategory( + HoodieWriteConflictException.ConflictCategory.TABLE_SERVICE_VS_INGESTION)); + assertDoesNotThrow(() -> metricsOff.emitConflictResolutionByCategory( + HoodieWriteConflictException.ConflictCategory.TABLE_SERVICE_VS_TABLE_SERVICE)); + } + + // ----------------------------------------------------------------------- + // Conflict resolution counters and timer + // ----------------------------------------------------------------------- + + @Test + public void testConflictResolutionSuccessAndFailureCounters() { + when(writeConfig.isLockingMetricsEnabled()).thenReturn(true); + + String successName = hoodieMetrics.getMetricsName(HoodieMetrics.CONFLICT_RESOLUTION_STR, HoodieMetrics.SUCCESS_COUNTER); + String failureName = hoodieMetrics.getMetricsName(HoodieMetrics.CONFLICT_RESOLUTION_STR, HoodieMetrics.FAILURE_COUNTER); + + hoodieMetrics.emitConflictResolutionSuccessful(); + hoodieMetrics.emitConflictResolutionSuccessful(); + assertEquals(2, metrics.getRegistry().getCounters().get(successName).getCount()); + + hoodieMetrics.emitConflictResolutionFailed(); + assertEquals(1, metrics.getRegistry().getCounters().get(failureName).getCount()); + } + + @Test + public void testConflictResolutionTimerCtx() throws InterruptedException { + when(writeConfig.isLockingMetricsEnabled()).thenReturn(true); + + Timer.Context ctx = hoodieMetrics.getConflictResolutionCtx(); + assertNotNull(ctx); + Thread.sleep(5); + assertTrue(hoodieMetrics.getDurationInMs(ctx.stop()) > 0); + } + + // ----------------------------------------------------------------------- + // Compaction counters + // ----------------------------------------------------------------------- + + @Test + public void testCompactionCounters() { + String requestedName = hoodieMetrics.getMetricsName( + HoodieTimeline.COMPACTION_ACTION, + HoodieTimeline.REQUESTED_COMPACTION_SUFFIX + HoodieMetrics.COUNTER_METRIC_EXTENSION); + String completedName = hoodieMetrics.getMetricsName( + HoodieTimeline.COMPACTION_ACTION, + HoodieTimeline.COMPLETED_COMPACTION_SUFFIX + HoodieMetrics.COUNTER_METRIC_EXTENSION); + + hoodieMetrics.emitCompactionRequested(); + hoodieMetrics.emitCompactionRequested(); + assertEquals(2, metrics.getRegistry().getCounters().get(requestedName).getCount()); + + hoodieMetrics.emitCompactionCompleted(); + assertEquals(1, metrics.getRegistry().getCounters().get(completedName).getCount()); + } + + // ----------------------------------------------------------------------- + // Archive metrics + // ----------------------------------------------------------------------- + + @Test + public void testArchiveTimerAndMetrics() throws InterruptedException { + Timer.Context ctx = hoodieMetrics.getArchiveCtx(); + assertNotNull(ctx); + Thread.sleep(5); + int numInstantsArchived = 7; + hoodieMetrics.updateArchiveMetrics(hoodieMetrics.getDurationInMs(ctx.stop()), numInstantsArchived); + + String durationName = hoodieMetrics.getMetricsName(HoodieMetrics.ARCHIVE_ACTION, HoodieMetrics.DURATION_STR); + String countName = hoodieMetrics.getMetricsName(HoodieMetrics.ARCHIVE_ACTION, HoodieMetrics.DELETE_INSTANTS_NUM_STR); + assertTrue((Long) metrics.getRegistry().getGauges().get(durationName).getValue() > 0); + assertEquals(numInstantsArchived, (long) metrics.getRegistry().getGauges().get(countName).getValue()); + } + + @Test + public void testUpdateArchivalMetrics() { + Map<String, Long> archivalMetrics = new HashMap<>(); + archivalMetrics.put("numFilesArchived", 10L); + archivalMetrics.put("archiveDurationMs", 250L); + + hoodieMetrics.updateArchivalMetrics(archivalMetrics); + + assertEquals(10L, (long) metrics.getRegistry().getGauges().get( + hoodieMetrics.getMetricsName("archival", "numFilesArchived")).getValue()); + assertEquals(250L, (long) metrics.getRegistry().getGauges().get( + hoodieMetrics.getMetricsName("archival", "archiveDurationMs")).getValue()); + } + + // ----------------------------------------------------------------------- + // Post-commit metrics + // ----------------------------------------------------------------------- + + @Test + public void testPostCommitMetrics() { + String successName = hoodieMetrics.getMetricsName(HoodieMetrics.POST_COMMIT_STR, HoodieMetrics.SUCCESS_COUNTER); + String failureName = hoodieMetrics.getMetricsName(HoodieMetrics.POST_COMMIT_STR, HoodieMetrics.FAILURE_COUNTER); + String durationName = hoodieMetrics.getMetricsName(HoodieMetrics.POST_COMMIT_STR, HoodieMetrics.DURATION_STR); + + hoodieMetrics.updatePostCommitMetrics(true, 100L); + hoodieMetrics.updatePostCommitMetrics(true, 200L); + hoodieMetrics.updatePostCommitMetrics(false, 50L); + + assertEquals(2, metrics.getRegistry().getCounters().get(successName).getCount()); + assertEquals(1, metrics.getRegistry().getCounters().get(failureName).getCount()); + assertEquals(50L, (long) metrics.getRegistry().getGauges().get(durationName).getValue()); + } + + // ----------------------------------------------------------------------- + // reportMetrics / updateClusteringFileCreationMetrics + // ----------------------------------------------------------------------- + + @Test + public void testReportMetrics() { + hoodieMetrics.reportMetrics("commit", "customMetric", 42L); + String metricName = hoodieMetrics.getMetricsName("commit", "customMetric"); + assertEquals(42L, (long) metrics.getRegistry().getGauges().get(metricName).getValue()); + } + + @Test + public void testUpdateClusteringFileCreationMetrics() { + hoodieMetrics.updateClusteringFileCreationMetrics(150L); + String metricName = hoodieMetrics.getMetricsName(HoodieTimeline.CLUSTERING_ACTION, "fileCreationTime"); + assertEquals(150L, (long) metrics.getRegistry().getGauges().get(metricName).getValue()); + } + + // ----------------------------------------------------------------------- + // Log-compaction and clustering timer contexts + // ----------------------------------------------------------------------- + + @Test + public void testLogCompactionTimerCtx() throws InterruptedException { + Timer.Context ctx = hoodieMetrics.getLogCompactionCtx(); + assertNotNull(ctx); + Thread.sleep(5); + assertTrue(hoodieMetrics.getDurationInMs(ctx.stop()) > 0); + } + + @Test + public void testClusteringTimerCtx() throws InterruptedException { + Timer.Context ctx = hoodieMetrics.getClusteringCtx(); + assertNotNull(ctx); + Thread.sleep(5); + assertTrue(hoodieMetrics.getDurationInMs(ctx.stop()) > 0); + } + + @Test + public void testCleanTimerCtx() throws InterruptedException { + Timer.Context ctx = hoodieMetrics.getCleanCtx(); + assertNotNull(ctx); + Thread.sleep(5); + assertTrue(hoodieMetrics.getDurationInMs(ctx.stop()) > 0); + } + + // ----------------------------------------------------------------------- + // Version metrics + // ----------------------------------------------------------------------- + + @Test + public void testVersionMetrics() { + hoodieMetrics.emitVersionMetrics(); + assertTrue(metrics.getRegistry().getGauges().keySet().stream() + .anyMatch(name -> name.startsWith("version."))); + } + + // ----------------------------------------------------------------------- + // Commit metrics with event time (latency / freshness) + // ----------------------------------------------------------------------- + + @Test + public void testCommitMetricsWithEventTime() { + long commitEpochTimeMs = System.currentTimeMillis(); + long durationMs = 1000L; + long minEventTimeMs = commitEpochTimeMs - 5000L; + long maxEventTimeMs = commitEpochTimeMs - 2000L; + + HoodieCommitMetadata metadata = mock(HoodieCommitMetadata.class); + when(metadata.fetchTotalPartitionsWritten()).thenReturn(1L); + when(metadata.fetchTotalFilesInsert()).thenReturn(0L); + when(metadata.fetchTotalFilesUpdated()).thenReturn(0L); + when(metadata.fetchTotalRecordsWritten()).thenReturn(10L); + when(metadata.fetchTotalUpdateRecordsWritten()).thenReturn(5L); + when(metadata.fetchTotalInsertRecordsWritten()).thenReturn(5L); + when(metadata.fetchTotalBytesWritten()).thenReturn(1024L); + when(metadata.getTotalScanTime()).thenReturn(0L); + when(metadata.getTotalCreateTime()).thenReturn(0L); + when(metadata.getTotalUpsertTime()).thenReturn(0L); + when(metadata.getTotalCompactedRecordsUpdated()).thenReturn(0L); + when(metadata.getTotalLogFilesCompacted()).thenReturn(0L); + when(metadata.getTotalLogFilesSize()).thenReturn(0L); + when(metadata.getTotalRecordsDeleted()).thenReturn(0L); + when(metadata.getMinAndMaxEventTime()).thenReturn(Pair.of(Option.of(minEventTimeMs), Option.of(maxEventTimeMs))); + when(writeConfig.isCompactionLogBlockMetricsOn()).thenReturn(false); + + hoodieMetrics.updateCommitMetrics(commitEpochTimeMs, durationMs, metadata, "commit"); + + long expectedLatency = commitEpochTimeMs + durationMs - minEventTimeMs; + long expectedFreshness = commitEpochTimeMs + durationMs - maxEventTimeMs; + assertEquals(expectedLatency, (long) metrics.getRegistry().getGauges().get( + hoodieMetrics.getMetricsName("commit", HoodieMetrics.COMMIT_LATENCY_IN_MS_STR)).getValue()); + assertEquals(expectedFreshness, (long) metrics.getRegistry().getGauges().get( + hoodieMetrics.getMetricsName("commit", HoodieMetrics.COMMIT_FRESHNESS_IN_MS_STR)).getValue()); + } + + // ----------------------------------------------------------------------- + // Rollback failure with null exception type + // ----------------------------------------------------------------------- + + @Test + public void testEmitRollbackFailureWithNullExceptionType() { + hoodieMetrics.emitRollbackFailure(null); + + String failureName = hoodieMetrics.getMetricsName("rollback", FAILURE_COUNTER); + assertEquals(1, metrics.getRegistry().getCounters().get(failureName).getCount()); + // No per-exception counter should be registered + long exceptionCounters = metrics.getRegistry().getCounters().keySet().stream() + .filter(n -> n.startsWith(hoodieMetrics.getMetricsName("rollback", "")) && !n.equals(failureName)) + .count(); + assertEquals(0, exceptionCounters); + } + + // ----------------------------------------------------------------------- + // getMetricsName prefix handling + // ----------------------------------------------------------------------- + + @Test + public void testGetMetricsNameWithPrefix() { + when(writeConfig.getMetricReporterMetricsNamePrefix()).thenReturn("my_prefix"); + assertEquals("my_prefix.action.metric", hoodieMetrics.getMetricsName("action", "metric")); + } + + @Test + public void testGetMetricsNameWithoutPrefix() { + when(writeConfig.getMetricReporterMetricsNamePrefix()).thenReturn(""); + assertEquals("action.metric", hoodieMetrics.getMetricsName("action", "metric")); + } + + // ----------------------------------------------------------------------- + // Existing rollback-failure and conflict-resolution-by-category tests + // ----------------------------------------------------------------------- + @Test public void testRollbackFailureMetric() { // Test that rollback failure metric is emitted correctly diff --git a/hudi-sync/hudi-sync-common/src/main/java/org/apache/hudi/sync/common/metrics/HoodieMetaSyncMetrics.java b/hudi-sync/hudi-sync-common/src/main/java/org/apache/hudi/sync/common/metrics/HoodieMetaSyncMetrics.java index c21adf7bc812..c0ea52b47575 100644 --- a/hudi-sync/hudi-sync-common/src/main/java/org/apache/hudi/sync/common/metrics/HoodieMetaSyncMetrics.java +++ b/hudi-sync/hudi-sync-common/src/main/java/org/apache/hudi/sync/common/metrics/HoodieMetaSyncMetrics.java @@ -78,6 +78,9 @@ public class HoodieMetaSyncMetrics { } public void incrementRecreateAndSyncFailureCounter() { + if (!metricsConfig.isMetricsOn()) { + return; + } recreateAndSyncFailureCounter = getCounter(recreateAndSyncFailureCounter, recreateAndSyncFailureCounterName); recreateAndSyncFailureCounter.inc(); }
