This is an automated email from the ASF dual-hosted git repository.

codope pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new aadeeaec709 [HUDI-8362] Fix DebeziumSource to always returning 
dataframe with schema (#12066)
aadeeaec709 is described below

commit aadeeaec709b35f88270547ad5d2bd99dd4f0b35
Author: fahmiduldul <[email protected]>
AuthorDate: Mon Dec 9 21:22:25 2024 +0700

    [HUDI-8362] Fix DebeziumSource to always returning dataframe with schema 
(#12066)
    
    * fix: DebeziumSource always returning dataframe with schema
    
    * feat: add unit test
    
    * fix: merge conflict
---
 .../utilities/sources/debezium/DebeziumSource.java | 28 ++++++++--------------
 .../debezium/TestAbstractDebeziumSource.java       | 18 ++++++++++++++
 2 files changed, 28 insertions(+), 18 deletions(-)

diff --git 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/debezium/DebeziumSource.java
 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/debezium/DebeziumSource.java
index 0a92dc1e684..e4dc95e26ec 100644
--- 
a/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/debezium/DebeziumSource.java
+++ 
b/hudi-utilities/src/main/java/org/apache/hudi/utilities/sources/debezium/DebeziumSource.java
@@ -117,24 +117,16 @@ public abstract class DebeziumSource extends RowSource {
     long totalNewMsgs = CheckpointUtils.totalNewMessages(offsetRanges);
     LOG.info("About to read " + totalNewMsgs + " from Kafka for topic :" + 
offsetGen.getTopicName());
 
-    if (totalNewMsgs == 0) {
-      // If there are no new messages, use empty dataframe with no schema. 
This is because the schema from schema registry can only be considered
-      // up to date if a change event has occurred.
-      return Pair.of(Option.of(sparkSession.emptyDataFrame()),
-          new StreamerCheckpointV2(overrideCheckpointStr.isEmpty()
-              ? CheckpointUtils.offsetsToStr(offsetRanges) : 
overrideCheckpointStr));
-    } else {
-      try {
-        String schemaStr = 
schemaRegistryProvider.fetchSchemaFromRegistry(getStringWithAltKeys(props, 
HoodieSchemaProviderConfig.SRC_SCHEMA_REGISTRY_URL));
-        Dataset<Row> dataset = toDataset(offsetRanges, offsetGen, schemaStr);
-        LOG.info(String.format("Spark schema of Kafka Payload for topic 
%s:\n%s", offsetGen.getTopicName(), dataset.schema().treeString()));
-        LOG.info(String.format("New checkpoint string: %s", 
CheckpointUtils.offsetsToStr(offsetRanges)));
-        return Pair.of(Option.of(dataset),
-            new StreamerCheckpointV2(overrideCheckpointStr.isEmpty() ? 
CheckpointUtils.offsetsToStr(offsetRanges) : overrideCheckpointStr));
-      } catch (Exception e) {
-        LOG.error("Fatal error reading and parsing incoming debezium event", 
e);
-        throw new HoodieReadFromSourceException("Fatal error reading and 
parsing incoming debezium event", e);
-      }
+    try {
+      String schemaStr = 
schemaRegistryProvider.fetchSchemaFromRegistry(getStringWithAltKeys(props, 
HoodieSchemaProviderConfig.SRC_SCHEMA_REGISTRY_URL));
+      Dataset<Row> dataset = toDataset(offsetRanges, offsetGen, schemaStr);
+      LOG.info(String.format("Spark schema of Kafka Payload for topic 
%s:\n%s", offsetGen.getTopicName(), dataset.schema().treeString()));
+      LOG.info(String.format("New checkpoint string: %s", 
CheckpointUtils.offsetsToStr(offsetRanges)));
+      return Pair.of(Option.of(dataset),
+              new StreamerCheckpointV2(overrideCheckpointStr.isEmpty() ? 
CheckpointUtils.offsetsToStr(offsetRanges) : overrideCheckpointStr));
+    } catch (Exception e) {
+      LOG.error("Fatal error reading and parsing incoming debezium event", e);
+      throw new HoodieReadFromSourceException("Fatal error reading and parsing 
incoming debezium event", e);
     }
   }
 
diff --git 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/debezium/TestAbstractDebeziumSource.java
 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/debezium/TestAbstractDebeziumSource.java
index 9e5d3d1f132..1865c8ffa60 100644
--- 
a/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/debezium/TestAbstractDebeziumSource.java
+++ 
b/hudi-utilities/src/test/java/org/apache/hudi/utilities/sources/debezium/TestAbstractDebeziumSource.java
@@ -39,6 +39,7 @@ import org.apache.spark.api.java.JavaSparkContext;
 import org.apache.spark.sql.Dataset;
 import org.apache.spark.sql.Row;
 import org.apache.spark.streaming.kafka010.KafkaTestUtils;
+import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeAll;
@@ -138,6 +139,23 @@ public abstract class TestAbstractDebeziumSource extends 
UtilitiesTestBase {
     validateMetaFields(fetch.getBatch().get());
   }
 
+  @Test
+  public void testDatasetRowSchemaWithoutData() throws Exception {
+    String sourceClass = getSourceClass();
+
+    // topic setup without message
+    testUtils.createTopic(testTopicName, 2);
+    TypedProperties props = createPropsForJsonSource();
+
+    SchemaProvider schemaProvider = new MockSchemaRegistryProvider(props, jsc, 
this);
+    SourceFormatAdapter debeziumSource = new 
SourceFormatAdapter(UtilHelpers.createSource(sourceClass, props, jsc, 
sparkSession, metrics, new DefaultStreamContext(schemaProvider, 
Option.empty())));
+    InputBatch<Dataset<Row>> fetch = 
debeziumSource.fetchNewDataInRowFormat(Option.empty(), 10);
+    Dataset<Row> result = fetch.getBatch().get();
+
+    assertEquals(result.count(), 0);
+    assertTrue(result.columns().length > 0);
+  }
+
   private GenericRecord generateDebeziumEvent(Operation op) {
     Schema schema = new Schema.Parser().parse(getSchema());
     String indexName = getIndexName().concat(".ghschema.gharchive.Value");

Reply via email to