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 cce4bbebed56ca805f9750406e4e844da06908e8
Author: Prashant Wason <[email protected]>
AuthorDate: Sat Jun 13 07:12:44 2026 +0800

    feat: add more metrics for delta streamer (#18085)
    
    * Add more metrics for hudi streamer
    
    Summary:
    Add the following metrics
    1. Emit success/failure metrics when streamer job is finished.
    2. Emit the number of commits to process in the streamer job.
    3. Emit the number of unprocessed commits in the source table timeline.
    
    ---------
    
    Co-authored-by: Jintao Guan <[email protected]>
    Co-authored-by: Claude Opus 4.6 <[email protected]>
    Co-authored-by: sivabalan <[email protected]>
    Co-authored-by: Lokesh Jain <[email protected]>
    
    (cherry picked from commit a032431c81aa71d5d627b13cb5189d9a9c54ba46)
---
 .../ingestion/HoodieIngestionMetrics.java          |  4 ++
 .../hudi/utilities/sources/SqlFileBasedSource.java |  6 ++-
 .../apache/hudi/utilities/sources/SqlSource.java   | 10 +++--
 .../hudi/utilities/streamer/HoodieStreamer.java    |  4 +-
 .../utilities/streamer/HoodieStreamerMetrics.java  | 15 +++++++
 .../apache/hudi/utilities/streamer/StreamSync.java |  8 ++++
 .../deltastreamer/TestHoodieDeltaStreamer.java     | 16 +++++---
 ...estHoodieDeltaStreamerSchemaEvolutionQuick.java |  9 ++++-
 .../TestHoodieDeltaStreamerWithMultiWriter.java    |  5 ++-
 .../multisync/TestMultipleMetaSync.java            | 13 +++---
 .../utilities/sources/TestSqlFileBasedSource.java  | 13 +++---
 .../hudi/utilities/sources/TestSqlSource.java      | 15 ++++---
 .../TestStreamerSourceCheckpointVersion.java       |  2 +-
 .../TestHoodieIncrSourceE2EAutoUpgrade.java        |  7 +++-
 .../streamer/TestHoodieStreamerMetrics.java        | 47 ++++++++++++++++++++++
 .../hudi/utilities/streamer/TestStreamSync.java    | 33 +++++++++++++++
 16 files changed, 173 insertions(+), 34 deletions(-)

diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/ingestion/HoodieIngestionMetrics.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/ingestion/HoodieIngestionMetrics.java
index d58e7877f5e3..44871f8c45fb 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/ingestion/HoodieIngestionMetrics.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/ingestion/HoodieIngestionMetrics.java
@@ -48,6 +48,10 @@ public abstract class HoodieIngestionMetrics {
 
   public abstract void updateStreamerMetrics(long durationNanos);
 
+  public abstract void emitStreamerJobSuccessMetrics();
+
+  public abstract void emitStreamerJobFailedMetrics();
+
   public abstract void updateStreamerMetaSyncMetrics(String 
syncClassShortName, long syncTimeNanos);
 
   public abstract void updateStreamerSyncMetrics(long syncEpochTimeMs);
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlFileBasedSource.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlFileBasedSource.java
index 913806420323..5ca569136e7c 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlFileBasedSource.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlFileBasedSource.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.Pair;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.hadoop.fs.HadoopFSUtils;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics;
 import org.apache.hudi.utilities.schema.SchemaProvider;
 
 import lombok.extern.slf4j.Slf4j;
@@ -65,16 +66,19 @@ public class SqlFileBasedSource extends RowSource {
 
   private final String sourceSqlFile;
   private final boolean shouldEmitCheckPoint;
+  private HoodieIngestionMetrics metrics;
 
   public SqlFileBasedSource(
       TypedProperties props,
       JavaSparkContext sparkContext,
       SparkSession sparkSession,
-      SchemaProvider schemaProvider) {
+      SchemaProvider schemaProvider,
+      HoodieIngestionMetrics metrics) {
     super(props, sparkContext, sparkSession, schemaProvider);
     checkRequiredConfigProperties(props, 
Collections.singletonList(SOURCE_SQL_FILE));
     sourceSqlFile = getStringWithAltKeys(props, SOURCE_SQL_FILE);
     shouldEmitCheckPoint = getBooleanWithAltKeys(props, EMIT_EPOCH_CHECKPOINT);
+    this.metrics = metrics;
   }
 
   @Override
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlSource.java 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlSource.java
index a345e937003f..21f1230db469 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlSource.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/SqlSource.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.table.checkpoint.Checkpoint;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.Pair;
 import org.apache.hudi.utilities.config.SqlSourceConfig;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics;
 import org.apache.hudi.utilities.schema.SchemaProvider;
 
 import lombok.extern.slf4j.Slf4j;
@@ -61,17 +62,20 @@ public class SqlSource extends RowSource {
   private static final long serialVersionUID = 1L;
   private final String sourceSql;
   private final SparkSession spark;
+  private final HoodieIngestionMetrics metrics;
 
   public SqlSource(
       TypedProperties props,
       JavaSparkContext sparkContext,
       SparkSession sparkSession,
-      SchemaProvider schemaProvider) {
+      SchemaProvider schemaProvider,
+      HoodieIngestionMetrics metrics) {
     super(props, sparkContext, sparkSession, schemaProvider);
     checkRequiredConfigProperties(
         props, Collections.singletonList(SqlSourceConfig.SOURCE_SQL));
-    sourceSql = getStringWithAltKeys(props, SqlSourceConfig.SOURCE_SQL);
-    spark = sparkSession;
+    this.sourceSql = getStringWithAltKeys(props, SqlSourceConfig.SOURCE_SQL);
+    this.spark = sparkSession;
+    this.metrics = metrics;
   }
 
   @Override
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
index 0186e5fd0584..9fc9dbff66a9 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamer.java
@@ -924,7 +924,9 @@ public class HoodieStreamer implements Serializable {
     public void ingestOnce() {
       try {
         streamSync.syncOnce();
-      } catch (IOException e) {
+        streamSync.reportSuccessMetrics();
+      } catch (Exception e) {
+        streamSync.reportFailureMetrics();
         throw new HoodieIngestionException(String.format("Ingestion via %s 
failed with exception.", this.getClass()), e);
       } finally {
         close();
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerMetrics.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerMetrics.java
index 5813533e2184..0e0db770012a 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerMetrics.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/HoodieStreamerMetrics.java
@@ -108,6 +108,20 @@ public class HoodieStreamerMetrics extends 
HoodieIngestionMetrics {
     }
   }
 
+  @Override
+  public void emitStreamerJobSuccessMetrics() {
+    if (writeConfig.isMetricsOn()) {
+      metrics.registerGauge(getMetricsName("deltastreamer", "success"), 1);
+    }
+  }
+
+  @Override
+  public void emitStreamerJobFailedMetrics() {
+    if (writeConfig.isMetricsOn()) {
+      metrics.registerGauge(getMetricsName("deltastreamer", "failure"), 1);
+    }
+  }
+
   @Override
   public void updateStreamerMetaSyncMetrics(String syncClassShortName, long 
syncNs) {
     if (writeConfig.isMetricsOn()) {
@@ -198,6 +212,7 @@ public class HoodieStreamerMetrics extends 
HoodieIngestionMetrics {
     }
   }
 
+  @Override
   public void updateStreamerSourceBytesToBeIngestedInSyncRound(long 
sourceBytesToBeIngested) {
     if (writeConfig.isMetricsOn()) {
       metrics.registerGauge(getMetricsName("deltastreamer", 
"sourceBytesToBeIngestedInSyncRound"), sourceBytesToBeIngested);
diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java
index 7f38d80f8308..e79dbece90b7 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/streamer/StreamSync.java
@@ -589,6 +589,14 @@ public class StreamSync implements Serializable, Closeable 
{
     }
   }
 
+  public void reportSuccessMetrics() {
+    metrics.emitStreamerJobSuccessMetrics();
+  }
+
+  public void reportFailureMetrics() {
+    metrics.emitStreamerJobFailedMetrics();
+  }
+
   private Option<String> 
getLastPendingClusteringInstant(Option<HoodieTimeline> commitTimelineOpt) {
     if (commitTimelineOpt.isPresent()) {
       Option<HoodieInstant> pendingClusteringInstant = 
commitTimelineOpt.get().getLastPendingClusterInstant();
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
index 0584a1ba2675..684352d03309 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamer.java
@@ -103,6 +103,7 @@ import 
org.apache.hudi.utilities.HoodieMetadataTableValidator;
 import org.apache.hudi.utilities.UtilHelpers;
 import org.apache.hudi.utilities.config.HoodieStreamerConfig;
 import org.apache.hudi.utilities.config.SourceTestConfig;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionException;
 import org.apache.hudi.utilities.schema.FilebasedSchemaProvider;
 import org.apache.hudi.utilities.schema.KafkaOffsetPostProcessor;
 import org.apache.hudi.utilities.schema.SchemaProvider;
@@ -2608,11 +2609,12 @@ public class TestHoodieDeltaStreamer extends 
HoodieDeltaStreamerTestBase {
     HoodieDeltaStreamer.Config cfg = TestHelpers.makeConfig(tableBasePath, 
WriteOperationType.BULK_INSERT,
         Collections.singletonList(SqlQueryBasedTransformer.class.getName()), 
PROPS_FILENAME_TEST_SOURCE, true,
         false, false, null, null);
-    Exception e = assertThrows(HoodieException.class, () -> {
+    Exception e = assertThrows(HoodieIngestionException.class, () -> {
       syncOnce(new HoodieDeltaStreamer(cfg, jsc, fs, 
hiveServer.getHiveConf()));
     }, "Should error out when schema provider is not provided");
     log.debug("Expected error during reading data from source ", e);
-    assertTrue(e.getMessage().contains("Schema provider is required for this 
operation and for the source of interest. "
+    String errorMsg = e.getCause() != null ? e.getCause().getMessage() : 
e.getMessage();
+    assertTrue(errorMsg.contains("Schema provider is required for this 
operation and for the source of interest. "
         + "Please set '--schemaprovider-class' in the top level HoodieStreamer 
config for the source of interest. "
         + "Based on the schema provider class chosen, additional configs might 
be required. "
         + "For eg, if you choose 
'org.apache.hudi.utilities.schema.SchemaRegistryProvider', "
@@ -3532,16 +3534,18 @@ public class TestHoodieDeltaStreamer extends 
HoodieDeltaStreamerTestBase {
     // Target schema is determined based on the Dataframe after transformation
     // No CSV header and no schema provider at the same time are not 
recommended,
     // as the transformer behavior may be unexpected
-    Exception e = assertThrows(AnalysisException.class, () -> {
+    Exception e = assertThrows(HoodieIngestionException.class, () -> {
       testCsvDFSSource(false, '\t', false, 
Collections.singletonList(TripsWithDistanceTransformer.class.getName()));
     }, "Should error out when doing the transformation.");
     log.debug("Expected error during transformation", e);
+    Throwable cause = e.getCause();
+    assertTrue(cause instanceof AnalysisException, "Expected cause to be 
AnalysisException but was: " + cause.getClass());
     // First message for Spark 3.4 and above, second message for Spark 3.3, 
third message for Spark 3.2 and below
     assertTrue(
-        e.getMessage().contains("[UNRESOLVED_COLUMN.WITH_SUGGESTION] A column 
or function parameter "
+        cause.getMessage().contains("[UNRESOLVED_COLUMN.WITH_SUGGESTION] A 
column or function parameter "
             + "with name `begin_lat` cannot be resolved. Did you mean one of 
the following?")
-            || e.getMessage().contains("Column 'begin_lat' does not exist. Did 
you mean one of the following?")
-            || e.getMessage().contains("cannot resolve 'begin_lat' given input 
columns:"));
+            || cause.getMessage().contains("Column 'begin_lat' does not exist. 
Did you mean one of the following?")
+            || cause.getMessage().contains("cannot resolve 'begin_lat' given 
input columns:"));
   }
 
   @Test
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerSchemaEvolutionQuick.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerSchemaEvolutionQuick.java
index e2fc16dc6a2f..5c07ff9d1294 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerSchemaEvolutionQuick.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerSchemaEvolutionQuick.java
@@ -29,6 +29,7 @@ import org.apache.hudi.common.table.timeline.HoodieInstant;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.exception.MissingSchemaFieldException;
 import org.apache.hudi.utilities.UtilHelpers;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionException;
 import org.apache.hudi.utilities.streamer.HoodieStreamer;
 
 import org.apache.spark.sql.Column;
@@ -231,7 +232,9 @@ public class TestHoodieDeltaStreamerSchemaEvolutionQuick 
extends TestHoodieDelta
       addData(df, false);
       deltaStreamer.sync();
       assertTrue(allowNullForDeletedCols);
-    } catch (MissingSchemaFieldException e) {
+    } catch (HoodieIngestionException e) {
+      assertTrue(e.getCause() instanceof MissingSchemaFieldException,
+          "Expected cause to be MissingSchemaFieldException but was: " + 
e.getCause());
       assertFalse(allowNullForDeletedCols);
       return;
     }
@@ -421,7 +424,9 @@ public class TestHoodieDeltaStreamerSchemaEvolutionQuick 
extends TestHoodieDelta
       assertTrue(riderFieldOpt.get().schema().getTypes()
           .stream().anyMatch(t -> HoodieSchemaType.STRING == t.getType()));
       
assertTrue(metaClient.reloadActiveTimeline().lastInstant().get().compareTo(lastInstant)
 > 0);
-    } catch (MissingSchemaFieldException e) {
+    } catch (HoodieIngestionException e) {
+      assertTrue(e.getCause() instanceof MissingSchemaFieldException,
+          "Expected cause to be MissingSchemaFieldException but was: " + 
e.getCause());
       assertFalse(allowNullForDeletedCols || targetSchemaSameAsTableSchema);
     }
   }
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerWithMultiWriter.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerWithMultiWriter.java
index bfe5dbda8190..194b9b9bb9d0 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerWithMultiWriter.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/TestHoodieDeltaStreamerWithMultiWriter.java
@@ -27,6 +27,7 @@ import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.table.HoodieTableMetaClient;
 import org.apache.hudi.common.table.checkpoint.CheckpointUtils;
 import org.apache.hudi.common.table.timeline.HoodieTimeline;
+import org.apache.hudi.common.testutils.JavaTestUtils;
 import org.apache.hudi.io.util.FileIOUtils;
 import org.apache.hudi.config.HoodieCleanConfig;
 import org.apache.hudi.config.HoodieCompactionConfig;
@@ -437,7 +438,7 @@ public class TestHoodieDeltaStreamerWithMultiWriter extends 
HoodieDeltaStreamerT
        * Need to perform getMessage().contains since the exception coming
        * from {@link 
org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer.DeltaSyncService} 
gets wrapped many times into RuntimeExceptions.
        */
-      if (expectConflict && backfillFailed.get() && 
e.getCause().getMessage().contains(ConcurrentModificationException.class.getName()))
 {
+      if (expectConflict && backfillFailed.get() && 
JavaTestUtils.checkNestedExceptionContains(e, 
ConcurrentModificationException.class.getName())) {
         // expected ConcurrentModificationException since ingestion & backfill 
will have overlapping writes
         if (!continuousFailed.get()) {
           // if backfill job failed, shutdown the continuous job.
@@ -447,7 +448,7 @@ public class TestHoodieDeltaStreamerWithMultiWriter extends 
HoodieDeltaStreamerT
           // both backfill and ingestion job cannot fail.
           throw new HoodieException("Both backfilling and ingestion job failed 
", e);
         }
-      } else if (expectConflict && continuousFailed.get() && 
e.getCause().getMessage().contains("Ingestion service was shut down with 
exception")) {
+      } else if (expectConflict && continuousFailed.get() && 
JavaTestUtils.checkNestedExceptionContains(e, "Ingestion service was shut down 
with exception")) {
         // incase of regular ingestion job failing, 
ConcurrentModificationException is not throw all the way.
         if (!backfillFailed.get()) {
           log.warn("Calling shutdown on backfill job since the 
ingstion/continuous job has failed for {}", jobId);
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/multisync/TestMultipleMetaSync.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/multisync/TestMultipleMetaSync.java
index d1e0f70b78f8..9c332de1feba 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/multisync/TestMultipleMetaSync.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/deltastreamer/multisync/TestMultipleMetaSync.java
@@ -22,6 +22,7 @@ import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.exception.HoodieMetaSyncException;
 import org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer;
 import org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamerTestBase;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionException;
 import org.apache.hudi.utilities.schema.FilebasedSchemaProvider;
 import org.apache.hudi.utilities.sources.TestDataSource;
 import org.apache.hudi.utilities.testutils.UtilitiesTestBase;
@@ -61,8 +62,9 @@ public class TestMultipleMetaSync extends 
HoodieDeltaStreamerTestBase {
     MockSyncTool1.syncSuccess = false;
     MockSyncTool2.syncSuccess = false;
     HoodieDeltaStreamer.Config cfg = getConfig(tableBasePath, syncClassNames);
-    Exception e = assertThrows(HoodieMetaSyncException.class, () -> 
syncOnce(new HoodieDeltaStreamer(cfg, jsc, fs, hiveServer.getHiveConf())));
-    
assertTrue(e.getMessage().contains(MockSyncToolException1.class.getName()));
+    Exception e = assertThrows(HoodieIngestionException.class, () -> 
syncOnce(new HoodieDeltaStreamer(cfg, jsc, fs, hiveServer.getHiveConf())));
+    assertTrue(e.getCause() instanceof HoodieMetaSyncException);
+    
assertTrue(e.getCause().getMessage().contains(MockSyncToolException1.class.getName()));
     assertTrue(MockSyncTool1.syncSuccess);
     assertTrue(MockSyncTool2.syncSuccess);
   }
@@ -73,9 +75,10 @@ public class TestMultipleMetaSync extends 
HoodieDeltaStreamerTestBase {
     MockSyncTool1.syncSuccess = false;
     MockSyncTool2.syncSuccess = false;
     HoodieDeltaStreamer.Config cfg = getConfig(tableBasePath, 
getSyncNames("MockSyncTool1", "MockSyncTool2", "MockSyncToolException1", 
"MockSyncToolException2"));
-    Exception e = assertThrows(HoodieMetaSyncException.class, () -> 
syncOnce(new HoodieDeltaStreamer(cfg, jsc, fs, hiveServer.getHiveConf())));
-    
assertTrue(e.getMessage().contains(MockSyncToolException1.class.getName()));
-    
assertTrue(e.getMessage().contains(MockSyncToolException2.class.getName()));
+    Exception e = assertThrows(HoodieIngestionException.class, () -> 
syncOnce(new HoodieDeltaStreamer(cfg, jsc, fs, hiveServer.getHiveConf())));
+    assertTrue(e.getCause() instanceof HoodieMetaSyncException);
+    
assertTrue(e.getCause().getMessage().contains(MockSyncToolException1.class.getName()));
+    
assertTrue(e.getCause().getMessage().contains(MockSyncToolException2.class.getName()));
     assertTrue(MockSyncTool1.syncSuccess);
     assertTrue(MockSyncTool2.syncSuccess);
   }
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlFileBasedSource.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlFileBasedSource.java
index f89552e62390..72984e7247f9 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlFileBasedSource.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlFileBasedSource.java
@@ -24,6 +24,7 @@ import org.apache.hudi.common.table.checkpoint.Checkpoint;
 import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
 import org.apache.hudi.common.util.Option;
 import org.apache.hudi.common.util.collection.Pair;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics;
 import org.apache.hudi.utilities.schema.FilebasedSchemaProvider;
 import org.apache.hudi.utilities.streamer.SourceFormatAdapter;
 import org.apache.hudi.utilities.testutils.UtilitiesTestBase;
@@ -45,6 +46,7 @@ import java.io.IOException;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
 
 /**
  * Test against {@link SqlSource}.
@@ -60,6 +62,7 @@ public class TestSqlFileBasedSource extends UtilitiesTestBase 
{
   private TypedProperties props;
   private SqlFileBasedSource sqlFileSource;
   private SourceFormatAdapter sourceFormatAdapter;
+  private final HoodieIngestionMetrics metrics = 
mock(HoodieIngestionMetrics.class);
 
   @BeforeAll
   public static void initClass() throws Exception {
@@ -114,7 +117,7 @@ public class TestSqlFileBasedSource extends 
UtilitiesTestBase {
         UtilitiesTestBase.basePath + "/sql-file-based-source.sql");
 
     props.setProperty(sqlFileSourceConfig, UtilitiesTestBase.basePath + 
"/sql-file-based-source.sql");
-    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider);
+    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider, metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlFileSource);
 
     // Test fetching Avro format
@@ -141,7 +144,7 @@ public class TestSqlFileBasedSource extends 
UtilitiesTestBase {
         UtilitiesTestBase.basePath + "/sql-file-based-source.sql");
 
     props.setProperty(sqlFileSourceConfig, UtilitiesTestBase.basePath + 
"/sql-file-based-source.sql");
-    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider);
+    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider, metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlFileSource);
 
     // Test fetching Row format
@@ -163,7 +166,7 @@ public class TestSqlFileBasedSource extends 
UtilitiesTestBase {
         UtilitiesTestBase.basePath + "/sql-file-based-source.sql");
 
     props.setProperty(sqlFileSourceConfig, UtilitiesTestBase.basePath + 
"/sql-file-based-source.sql");
-    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider);
+    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider, metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlFileSource);
 
     InputBatch<Dataset<Row>> fetch1AsRows =
@@ -184,7 +187,7 @@ public class TestSqlFileBasedSource extends 
UtilitiesTestBase {
         UtilitiesTestBase.basePath + 
"/sql-file-based-source-invalid-table.sql");
 
     props.setProperty(sqlFileSourceConfig, UtilitiesTestBase.basePath + 
"/sql-file-based-source-invalid-table.sql");
-    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider);
+    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider, metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlFileSource);
 
     assertThrows(
@@ -201,7 +204,7 @@ public class TestSqlFileBasedSource extends 
UtilitiesTestBase {
     props.setProperty(sqlFileSourceConfig, UtilitiesTestBase.basePath + 
"/sql-file-based-source.sql");
     props.setProperty(sqlFileSourceConfigEmitChkPointConf, "true");
 
-    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider);
+    sqlFileSource = new SqlFileBasedSource(props, jsc, sparkSession, 
schemaProvider, metrics);
     Pair<Option<Dataset<Row>>, Checkpoint> nextBatch = 
sqlFileSource.fetchNextBatch(Option.empty(), Long.MAX_VALUE);
 
     assertEquals(10000, nextBatch.getLeft().get().count());
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlSource.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlSource.java
index 146c5253d5f5..8a393e8b982e 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlSource.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestSqlSource.java
@@ -22,6 +22,7 @@ import org.apache.hudi.AvroConversionUtils;
 import org.apache.hudi.common.config.TypedProperties;
 import org.apache.hudi.common.testutils.HoodieTestDataGenerator;
 import org.apache.hudi.common.util.Option;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics;
 import org.apache.hudi.utilities.schema.FilebasedSchemaProvider;
 import org.apache.hudi.utilities.streamer.SourceFormatAdapter;
 import org.apache.hudi.utilities.testutils.UtilitiesTestBase;
@@ -43,6 +44,7 @@ import java.io.IOException;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
 
 /**
  * Test against {@link SqlSource}.
@@ -57,6 +59,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   private TypedProperties props;
   private SqlSource sqlSource;
   private SourceFormatAdapter sourceFormatAdapter;
+  private final HoodieIngestionMetrics metrics = 
mock(HoodieIngestionMetrics.class);
 
   @BeforeAll
   public static void initClass() throws Exception {
@@ -105,7 +108,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceAvroFormat() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from test_sql_table");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     // Test fetching Avro format
@@ -128,7 +131,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceRowFormat() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from test_sql_table");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     // Test fetching Row format
@@ -146,7 +149,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceCheckpoint() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from test_sql_table where 
1=0");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     InputBatch<Dataset<Row>> fetch1AsRows =
@@ -163,7 +166,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceMoreRecordsThanSourceLimit() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from test_sql_table");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     InputBatch<Dataset<Row>> fetch1AsRows =
@@ -180,7 +183,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceZeroRecord() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from test_sql_table where 
1=0");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     InputBatch<Dataset<Row>> fetch1AsRows =
@@ -197,7 +200,7 @@ public class TestSqlSource extends UtilitiesTestBase {
   @Test
   public void testSqlSourceInvalidTable() throws IOException {
     props.setProperty(sqlSourceConfig, "select * from not_exist_sql_table");
-    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider);
+    sqlSource = new SqlSource(props, jsc, sparkSession, schemaProvider, 
metrics);
     sourceFormatAdapter = new SourceFormatAdapter(sqlSource);
 
     assertThrows(
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestStreamerSourceCheckpointVersion.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestStreamerSourceCheckpointVersion.java
index 32480348a565..70ec4e708b9b 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestStreamerSourceCheckpointVersion.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/TestStreamerSourceCheckpointVersion.java
@@ -168,7 +168,7 @@ class TestStreamerSourceCheckpointVersion {
     TypedProperties props = propsWith(writeTableVersion);
     props.setProperty("hoodie.streamer.source.sql.file", sqlFile.toString());
     props.setProperty("hoodie.streamer.source.sql.checkpoint.emit", "true");
-    SqlFileBasedSource source = new SqlFileBasedSource(props, jsc, spark, 
null);
+    SqlFileBasedSource source = new SqlFileBasedSource(props, jsc, spark, 
null, null);
     Pair<Option<Dataset<Row>>, Checkpoint> result =
         invokeRowSourceFetch(source, makeInputCheckpoint(inputKind, "k"));
     assertV1(result.getRight());
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieIncrSourceE2EAutoUpgrade.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieIncrSourceE2EAutoUpgrade.java
index 661216053b23..3fd2068c9100 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieIncrSourceE2EAutoUpgrade.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieIncrSourceE2EAutoUpgrade.java
@@ -37,6 +37,7 @@ import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.hudi.testutils.HoodieClientTestUtils;
 import org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer;
 import org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamerTestBase;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionException;
 import org.apache.hudi.utilities.sources.MockGeneralHoodieIncrSource;
 import org.apache.hudi.utilities.sources.S3EventsHoodieIncrSourceHarness;
 
@@ -191,8 +192,10 @@ public class TestHoodieIncrSourceE2EAutoUpgrade extends 
S3EventsHoodieIncrSource
     HoodieDeltaStreamer.Config cfg = createConfig(basePath(), null, 
sourceClass);
     cfg.checkpoint = "overrideWhenAutoUpgradingWouldFail";
     ds = new HoodieDeltaStreamer(cfg, jsc, Option.of(props));
-    Exception ex = assertThrows(HoodieUpgradeDowngradeException.class, 
ds::sync);
-    assertTrue(ex.getMessage().contains("When upgrade/downgrade is happening, 
please avoid setting --checkpoint option and --ignore-checkpoint for your delta 
streamers."));
+    Exception ex = assertThrows(HoodieIngestionException.class, ds::sync);
+    Throwable cause = ex.getCause();
+    assertTrue(cause instanceof HoodieUpgradeDowngradeException, "Expected 
cause to be HoodieUpgradeDowngradeException but was: " + cause.getClass());
+    assertTrue(cause.getMessage().contains("When upgrade/downgrade is 
happening, please avoid setting --checkpoint option and --ignore-checkpoint for 
your delta streamers."));
     // No changes to the timeline / table config.
     metaClient.reloadActiveTimeline();
     assertEquals(metaClient.getActiveTimeline().lastInstant().get(), 
instantAfterFirstRound);
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerMetrics.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerMetrics.java
index 8a45c79e7348..4fb7db1cd634 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerMetrics.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestHoodieStreamerMetrics.java
@@ -68,4 +68,51 @@ public class TestHoodieStreamerMetrics {
     metrics.updateErrorTableCommitDuration(0L);
     assertNull(metrics.getMetrics());
   }
+
+  @Test
+  public void testEmitStreamerJobSuccessMetrics() {
+    HoodieMetricsConfig metricsConfig = HoodieMetricsConfig.newBuilder()
+        .on(true)
+        .withPath("/tmp/path3")
+        .withReporterType("INMEMORY")
+        .build();
+    HoodieStreamerMetrics metrics = new HoodieStreamerMetrics(
+        metricsConfig, HoodieStorageUtils.getStorage(getDefaultStorageConf()));
+    metrics.emitStreamerJobSuccessMetrics();
+    MetricRegistry registry = metrics.getMetrics().getRegistry();
+    assertEquals(1, registry.getGauges().size());
+    assertEquals(".deltastreamer.success", registry.getGauges().firstKey());
+    assertEquals(1L, 
registry.getGauges().get(".deltastreamer.success").getValue());
+  }
+
+  @Test
+  public void testEmitStreamerJobFailedMetrics() {
+    HoodieMetricsConfig metricsConfig = HoodieMetricsConfig.newBuilder()
+        .on(true)
+        .withPath("/tmp/path4")
+        .withReporterType("INMEMORY")
+        .build();
+    HoodieStreamerMetrics metrics = new HoodieStreamerMetrics(
+        metricsConfig, HoodieStorageUtils.getStorage(getDefaultStorageConf()));
+    metrics.emitStreamerJobFailedMetrics();
+    MetricRegistry registry = metrics.getMetrics().getRegistry();
+    assertEquals(1, registry.getGauges().size());
+    assertEquals(".deltastreamer.failure", registry.getGauges().firstKey());
+    assertEquals(1L, 
registry.getGauges().get(".deltastreamer.failure").getValue());
+  }
+
+  @Test
+  public void testEmitStreamerJobMetricsIfDisabled() {
+    HoodieMetricsConfig metricsConfig = HoodieMetricsConfig.newBuilder()
+        .on(false)
+        .withPath("/tmp/path5")
+        .withReporterType("INMEMORY")
+        .build();
+    HoodieStreamerMetrics metrics = new HoodieStreamerMetrics(
+        metricsConfig, HoodieStorageUtils.getStorage(getDefaultStorageConf()));
+    // Should not throw when metrics are disabled
+    metrics.emitStreamerJobSuccessMetrics();
+    metrics.emitStreamerJobFailedMetrics();
+    assertNull(metrics.getMetrics());
+  }
 }
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestStreamSync.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestStreamSync.java
index 4a851e7c0a06..42499d8c76d4 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestStreamSync.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/streamer/TestStreamSync.java
@@ -37,6 +37,7 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.hadoop.HoodieHadoopStorage;
 import org.apache.hudi.testutils.SparkClientFunctionalTestHarness;
+import org.apache.hudi.utilities.ingestion.HoodieIngestionMetrics;
 import org.apache.hudi.utilities.schema.SchemaProvider;
 import org.apache.hudi.utilities.sources.InputBatch;
 import org.apache.hudi.utilities.transform.Transformer;
@@ -56,6 +57,7 @@ import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
 
 import java.io.IOException;
+import java.lang.reflect.Field;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
@@ -76,6 +78,7 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.doCallRealMethod;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
@@ -473,5 +476,35 @@ public class TestStreamSync extends 
SparkClientFunctionalTestHarness {
       assertEquals("org.apache.hudi.common.model.OverwriteWithLatestPayload", 
triple.getMiddle());
       assertEquals("any_id", triple.getRight());
     }
+
+    @Test
+    void testReportSuccessMetricsDelegatesToMetrics() throws Exception {
+      StreamSync streamSync = mock(StreamSync.class);
+      HoodieIngestionMetrics mockMetrics = mock(HoodieIngestionMetrics.class);
+      Field metricsField = StreamSync.class.getDeclaredField("metrics");
+      metricsField.setAccessible(true);
+      metricsField.set(streamSync, mockMetrics);
+      doCallRealMethod().when(streamSync).reportSuccessMetrics();
+
+      streamSync.reportSuccessMetrics();
+
+      verify(mockMetrics, times(1)).emitStreamerJobSuccessMetrics();
+      verify(mockMetrics, never()).emitStreamerJobFailedMetrics();
+    }
+
+    @Test
+    void testReportFailureMetricsDelegatesToMetrics() throws Exception {
+      StreamSync streamSync = mock(StreamSync.class);
+      HoodieIngestionMetrics mockMetrics = mock(HoodieIngestionMetrics.class);
+      Field metricsField = StreamSync.class.getDeclaredField("metrics");
+      metricsField.setAccessible(true);
+      metricsField.set(streamSync, mockMetrics);
+      doCallRealMethod().when(streamSync).reportFailureMetrics();
+
+      streamSync.reportFailureMetrics();
+
+      verify(mockMetrics, times(1)).emitStreamerJobFailedMetrics();
+      verify(mockMetrics, never()).emitStreamerJobSuccessMetrics();
+    }
   }
 }

Reply via email to