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");