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(); + } } }
