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

danny0405 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 ed7327b32589 test(flink): improve streamer config and schema provider 
coverage (#19392)
ed7327b32589 is described below

commit ed7327b325894405a49fb9ab1f7730f851a51b10
Author: Danny Chan <[email protected]>
AuthorDate: Wed Jul 29 17:34:41 2026 +0800

    test(flink): improve streamer config and schema provider coverage (#19392)
---
 .../apache/hudi/streamer/FlinkStreamerConfig.java  |   4 +-
 .../java/org/apache/hudi/util/StreamerUtil.java    |   8 +-
 .../hudi/configuration/TestOptionsInference.java   |  98 +++++++++
 .../hudi/configuration/TestOptionsResolver.java    | 161 ++++++++++++++-
 .../hudi/schema/TestFilebasedSchemaProvider.java   |  92 +++++++++
 .../hudi/schema/TestSchemaRegistryProvider.java    | 127 ++++++++++++
 .../sink/clustering/TestFlinkClusteringConfig.java | 110 ++++++++++
 .../sink/compact/TestFlinkCompactionConfig.java    |  79 ++++++++
 .../hudi/streamer/TestFlinkStreamerConfig.java     | 146 ++++++++++++++
 .../hudi/streamer/TestHoodieFlinkStreamer.java     | 224 +++++++++++++++++++++
 10 files changed, 1046 insertions(+), 3 deletions(-)

diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/streamer/FlinkStreamerConfig.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/streamer/FlinkStreamerConfig.java
index 167ce428c66e..e19312ee548a 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/streamer/FlinkStreamerConfig.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/streamer/FlinkStreamerConfig.java
@@ -42,6 +42,7 @@ import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 import static 
org.apache.hudi.common.util.PartitionPathEncodeUtils.DEFAULT_PARTITION_PATH;
 import static org.apache.hudi.configuration.FlinkOptions.PARTITION_FORMAT_DAY;
@@ -433,7 +434,8 @@ public class FlinkStreamerConfig extends Configuration {
     conf.set(FlinkOptions.RECORD_MERGER_STRATEGY_ID, 
config.recordMergerStrategy);
     conf.set(FlinkOptions.PRE_COMBINE, config.preCombine);
     conf.set(FlinkOptions.RETRY_TIMES, 
Integer.parseInt(config.instantRetryTimes));
-    conf.set(FlinkOptions.RETRY_INTERVAL_MS, 
Long.parseLong(config.instantRetryInterval));
+    conf.set(FlinkOptions.RETRY_INTERVAL_MS,
+        
TimeUnit.SECONDS.toMillis(Long.parseLong(config.instantRetryInterval)));
     conf.set(FlinkOptions.IGNORE_FAILED, config.commitOnErrors);
     conf.set(FlinkOptions.RECORD_KEY_FIELD, config.recordKeyField);
     conf.set(FlinkOptions.PARTITION_PATH_FIELD, config.partitionPathField);
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java
index 61ab6c5db932..099da2c66eeb 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/util/StreamerUtil.java
@@ -135,7 +135,13 @@ public class StreamerUtil {
 
   public static TypedProperties getProps(FlinkStreamerConfig cfg) {
     if (cfg.propsFilePath.isEmpty()) {
-      return new TypedProperties();
+      TypedProperties properties = new TypedProperties();
+      cfg.configs.forEach(x -> {
+        String[] kv = x.split("=");
+        ValidationUtils.checkArgument(kv.length == 2);
+        properties.setProperty(kv[0], kv[1]);
+      });
+      return properties;
     }
     return readConfig(
         HadoopConfigurations.getHadoopConf(cfg),
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsInference.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsInference.java
index 5d13ae0e67ea..0bc69ae910b8 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsInference.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsInference.java
@@ -22,7 +22,12 @@ import org.apache.hudi.common.model.WriteConcurrencyMode;
 import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.util.ClientIds;
 
+import org.apache.flink.FlinkVersion;
+import org.apache.flink.api.common.RuntimeExecutionMode;
 import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.ExecutionOptions;
+import org.apache.flink.configuration.JobManagerOptions;
+import org.apache.flink.configuration.SchedulerExecutionMode;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
@@ -30,6 +35,9 @@ import java.io.File;
 
 import static org.hamcrest.CoreMatchers.is;
 import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 /**
  * Test cases for {@link OptionsInference}.
@@ -38,6 +46,83 @@ public class TestOptionsInference {
   @TempDir
   File tempFile;
 
+  @Test
+  void testSetupSourceAndSinkTasks() {
+    Configuration conf = new Configuration();
+
+    OptionsInference.setupSourceTasks(conf, 3);
+    OptionsInference.setupSinkTasks(conf, 4);
+
+    assertEquals(3, conf.get(FlinkOptions.READ_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.WRITE_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.COMPACTION_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.CLUSTERING_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.INDEX_WRITE_TASKS));
+
+    conf.set(FlinkOptions.READ_TASKS, 7);
+    conf.set(FlinkOptions.WRITE_TASKS, 8);
+    conf.set(FlinkOptions.BUCKET_ASSIGN_TASKS, 9);
+    conf.set(FlinkOptions.COMPACTION_TASKS, 10);
+    conf.set(FlinkOptions.CLUSTERING_TASKS, 11);
+    conf.set(FlinkOptions.INDEX_WRITE_TASKS, 12);
+
+    OptionsInference.setupSourceTasks(conf, 20);
+    OptionsInference.setupSinkTasks(conf, 20);
+
+    assertEquals(7, conf.get(FlinkOptions.READ_TASKS));
+    assertEquals(8, conf.get(FlinkOptions.WRITE_TASKS));
+    assertEquals(9, conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
+    assertEquals(10, conf.get(FlinkOptions.COMPACTION_TASKS));
+    assertEquals(11, conf.get(FlinkOptions.CLUSTERING_TASKS));
+    assertEquals(12, conf.get(FlinkOptions.INDEX_WRITE_TASKS));
+  }
+
+  @Test
+  void testSetupRuntimeConfigurations() {
+    Configuration conf = new Configuration();
+    conf.set(JobManagerOptions.SCHEDULER, 
JobManagerOptions.SchedulerType.AdaptiveBatch);
+    Configuration runtimeConf = new Configuration();
+    runtimeConf.set(ExecutionOptions.RUNTIME_MODE, RuntimeExecutionMode.BATCH);
+
+    OptionsInference.setupRuntimeConfigs(conf, runtimeConf);
+
+    if (FlinkVersion.current().toString().compareTo("2.0") >= 0) {
+      
assertTrue(conf.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+    } else {
+      
assertFalse(conf.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+    }
+
+    conf.set(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION, false);
+    runtimeConf.set(ExecutionOptions.RUNTIME_MODE, 
RuntimeExecutionMode.STREAMING);
+    OptionsInference.setupRuntimeConfigs(conf, runtimeConf);
+    assertFalse(conf.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+  }
+
+  @Test
+  void testSchedulerTypeResolutionThroughRuntimeSetup() {
+    Configuration runtimeConf = new Configuration();
+    runtimeConf.set(ExecutionOptions.RUNTIME_MODE, RuntimeExecutionMode.BATCH);
+    boolean isFlink2 = FlinkVersion.current().toString().compareTo("2.0") >= 0;
+
+    Configuration reactive = new Configuration();
+    reactive.set(JobManagerOptions.SCHEDULER_MODE, 
SchedulerExecutionMode.REACTIVE);
+    OptionsInference.setupRuntimeConfigs(reactive, runtimeConf);
+    assertEquals(isFlink2,
+        reactive.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+
+    Configuration adaptive = new Configuration();
+    adaptive.set(JobManagerOptions.SCHEDULER, 
JobManagerOptions.SchedulerType.Adaptive);
+    OptionsInference.setupRuntimeConfigs(adaptive, runtimeConf);
+    assertEquals(isFlink2,
+        adaptive.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+
+    Configuration defaultScheduler = new Configuration();
+    defaultScheduler.set(JobManagerOptions.SCHEDULER, 
JobManagerOptions.SchedulerType.Default);
+    OptionsInference.setupRuntimeConfigs(defaultScheduler, runtimeConf);
+    
assertFalse(defaultScheduler.get(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION));
+  }
+
   @Test
   void testSetupClientId() throws Exception {
     Configuration conf = getConf();
@@ -75,4 +160,17 @@ public class TestOptionsInference {
     conf.set(FlinkOptions.PATH, tempFile.getAbsolutePath());
     return conf;
   }
+
+  @Test
+  void testClientIdAndIndexSetupAreNoOpsWhenNotApplicable() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.PATH, tempFile.getAbsolutePath());
+    conf.set(FlinkOptions.INDEX_TYPE, "BLOOM");
+
+    OptionsInference.setupClientId(conf);
+    OptionsInference.setupIndexConfigs(conf);
+
+    assertFalse(conf.contains(FlinkOptions.WRITE_CLIENT_ID));
+    
assertFalse(conf.contains(FlinkOptions.BUCKET_INDEX_PARTITION_EXPRESSIONS));
+  }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java
index bd22f6641584..1cbeb4fb9033 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/configuration/TestOptionsResolver.java
@@ -18,28 +18,41 @@
 
 package org.apache.hudi.configuration;
 
+import 
org.apache.hudi.client.transaction.BucketIndexConcurrentFileWritesConflictResolutionStrategy;
+import 
org.apache.hudi.client.transaction.SimpleConcurrentFileWritesConflictResolutionStrategy;
+import org.apache.hudi.common.config.HoodieCommonConfig;
 import org.apache.hudi.common.config.HoodieMetadataConfig;
+import org.apache.hudi.common.model.DefaultHoodieRecordPayload;
 import org.apache.hudi.common.model.HoodieFailedWritesCleaningPolicy;
 import org.apache.hudi.common.model.HoodieTableType;
 import org.apache.hudi.common.model.WriteConcurrencyMode;
 import org.apache.hudi.common.model.WriteOperationType;
 import org.apache.hudi.common.table.HoodieTableConfig;
+import org.apache.hudi.common.table.cdc.HoodieCDCSupplementalLoggingMode;
+import 
org.apache.hudi.common.table.timeline.TimelineUtils.HollowCommitHandling;
 import org.apache.hudi.config.HoodieCleanConfig;
 import org.apache.hudi.config.HoodieIndexConfig;
 import org.apache.hudi.config.HoodieWriteConfig;
+import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.index.HoodieIndex;
+import org.apache.hudi.keygen.constant.KeyGeneratorOptions;
+import org.apache.hudi.sink.buffer.BufferMemoryType;
 import org.apache.hudi.utils.TestConfigurations;
 
+import org.apache.flink.api.common.functions.Partitioner;
 import org.apache.flink.configuration.Configuration;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.io.TempDir;
 
 import java.io.File;
+import java.util.Collections;
 
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 
 /**
@@ -69,7 +82,7 @@ public class TestOptionsResolver {
 
   @TempDir
   File tempFile;
-  
+
   @Test
   void testGetIndexType() {
     Configuration conf = getConf();
@@ -270,4 +283,150 @@ public class TestOptionsResolver {
     
conf.setString(HoodieMetadataConfig.GLOBAL_RECORD_LEVEL_INDEX_MAX_FILE_GROUP_COUNT_PROP.key(),
 "11");
     assertEquals(11, OptionsResolver.estimateFileGroupCountForRLI(conf));
   }
+
+  @Test
+  void testIncrementalJobGraphPredicate() {
+    Configuration conf = new Configuration();
+    assertFalse(OptionsResolver.isIncrementalJobGraph(conf));
+    conf.set(FlinkOptions.WRITE_INCREMENTAL_JOB_GRAPH_GENERATION, true);
+    assertTrue(OptionsResolver.isIncrementalJobGraph(conf));
+  }
+
+  @Test
+  void testTableTypePredicates() {
+    Configuration conf = new Configuration();
+    assertTrue(OptionsResolver.isCowTable(conf));
+    assertFalse(OptionsResolver.isMorTable(conf));
+    assertFalse(OptionsResolver.isMorTable(Collections.emptyMap()));
+    conf.set(FlinkOptions.TABLE_TYPE, 
HoodieTableType.MERGE_ON_READ.name().toLowerCase());
+    assertTrue(OptionsResolver.isMorTable(conf));
+    assertTrue(OptionsResolver.isMorTable(
+        Collections.singletonMap(FlinkOptions.TABLE_TYPE.key(), 
HoodieTableType.MERGE_ON_READ.name())));
+  }
+
+  @Test
+  void testOperationTypePredicates() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.OPERATION, WriteOperationType.INSERT.value());
+    assertTrue(OptionsResolver.isInsertOperation(conf));
+    conf.set(FlinkOptions.OPERATION, WriteOperationType.UPSERT.value());
+    assertTrue(OptionsResolver.isUpsertOperation(conf));
+    conf.set(FlinkOptions.OPERATION, WriteOperationType.BULK_INSERT.value());
+    assertTrue(OptionsResolver.isBulkInsertOperation(conf));
+  }
+
+  @Test
+  void testPayloadAndCompactionPredicates() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.PAYLOAD_CLASS_NAME, 
DefaultHoodieRecordPayload.class.getName());
+    assertTrue(OptionsResolver.isDefaultHoodieRecordPayloadClazz(conf));
+    conf.set(FlinkOptions.COMPACTION_TRIGGER_STRATEGY, 
FlinkOptions.TIME_ELAPSED.toUpperCase());
+    assertTrue(OptionsResolver.isDeltaTimeCompaction(conf));
+    conf.set(FlinkOptions.COMPACTION_TRIGGER_STRATEGY, 
FlinkOptions.NUM_COMMITS);
+    assertFalse(OptionsResolver.isDeltaTimeCompaction(conf));
+  }
+
+  @Test
+  void testReadCommitsLimit() {
+    Configuration conf = new Configuration();
+    assertEquals(-1, OptionsResolver.getReadCommitsLimit(conf));
+    conf.set(FlinkOptions.READ_COMMITS_LIMIT, 5);
+    assertEquals(5, OptionsResolver.getReadCommitsLimit(conf));
+  }
+
+  @Test
+  void testCdcOptions() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.SUPPLEMENTAL_LOGGING_MODE,
+        
HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER.name().toLowerCase());
+    assertEquals(HoodieCDCSupplementalLoggingMode.DATA_BEFORE_AFTER,
+        OptionsResolver.getCDCSupplementalLoggingMode(conf));
+
+    conf.set(FlinkOptions.READ_CDC_FROM_CHANGELOG, false);
+    assertFalse(OptionsResolver.readCDCFromChangelog(conf));
+  }
+
+  @Test
+  void testIndexKeyFields() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.RECORD_KEY_FIELD, "id,tenant");
+    assertEquals("id", OptionsResolver.getIndexKeyFields(conf).get(0));
+    assertEquals("tenant", OptionsResolver.getIndexKeyFields(conf).get(1));
+  }
+
+  @Test
+  void testSchemaAndTimestampOptions() {
+    Configuration conf = new Configuration();
+    conf.setString(HoodieCommonConfig.SCHEMA_EVOLUTION_ENABLE.key(), "true");
+    assertTrue(OptionsResolver.isSchemaEvolutionEnabled(conf));
+    
conf.setString(KeyGeneratorOptions.KEYGENERATOR_CONSISTENT_LOGICAL_TIMESTAMP_ENABLED.key(),
 "true");
+    assertTrue(OptionsResolver.isConsistentLogicalTimestampEnabled(conf));
+    
conf.setString(HoodieCommonConfig.INCREMENTAL_READ_HANDLE_HOLLOW_COMMIT.key(),
+        HollowCommitHandling.USE_TRANSITION_TIME.name());
+    assertTrue(OptionsResolver.isReadByTxnCompletionTime(conf));
+  }
+
+  @Test
+  void testWriteFlags() {
+    Configuration conf = new Configuration();
+    conf.setString(HoodieWriteConfig.ALLOW_EMPTY_COMMIT.key(), "true");
+    assertTrue(OptionsResolver.allowCommitOnEmptyBatch(conf));
+    conf.setString(HoodieWriteConfig.COMPLEX_KEYGEN_NEW_ENCODING.key(), 
"true");
+    assertTrue(OptionsResolver.useComplexKeygenNewEncoding(conf));
+  }
+
+  @Test
+  void testConcurrencyControlModes() {
+    Configuration conf = new Configuration();
+    conf.setString(HoodieWriteConfig.WRITE_CONCURRENCY_MODE.key(),
+        WriteConcurrencyMode.NON_BLOCKING_CONCURRENCY_CONTROL.name());
+    assertTrue(OptionsResolver.isNonBlockingConcurrencyControl(conf));
+    conf.setString(HoodieWriteConfig.WRITE_CONCURRENCY_MODE.key(),
+        
WriteConcurrencyMode.OPTIMISTIC_CONCURRENCY_CONTROL.name().toLowerCase());
+    assertTrue(OptionsResolver.isOptimisticConcurrencyControl(conf));
+  }
+
+  @Test
+  void testInsertPartitioner() {
+    Configuration conf = new Configuration();
+    assertFalse(OptionsResolver.getInsertPartitioner(conf).isPresent());
+    conf.set(FlinkOptions.INSERT_PARTITIONER_CLASS_NAME, 
TestPartitioner.class.getName());
+    assertTrue(OptionsResolver.getInsertPartitioner(conf).isPresent());
+    conf.set(FlinkOptions.INSERT_PARTITIONER_CLASS_NAME, 
String.class.getName());
+    assertThrows(HoodieException.class, () -> 
OptionsResolver.getInsertPartitioner(conf));
+  }
+
+  @Test
+  void testConflictResolutionStrategies() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BLOOM.name());
+    
assertInstanceOf(SimpleConcurrentFileWritesConflictResolutionStrategy.class,
+        OptionsResolver.getConflictResolutionStrategy(conf));
+    conf.set(FlinkOptions.INDEX_TYPE, HoodieIndex.IndexType.BUCKET.name());
+    
assertInstanceOf(BucketIndexConcurrentFileWritesConflictResolutionStrategy.class,
+        OptionsResolver.getConflictResolutionStrategy(conf));
+  }
+
+  @Test
+  void testWriteBufferSizingAndManagedMemory() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.WRITE_TASK_MAX_SIZE, 300D);
+    conf.set(FlinkOptions.WRITE_MERGE_MAX_MEMORY, 50);
+    assertEquals(150L * 1024 * 1024, 
OptionsResolver.getWriteBufferSizeInBytes(conf));
+    conf.set(FlinkOptions.WRITE_TASK_MAX_SIZE, 100D);
+    assertThrows(IllegalStateException.class, () -> 
OptionsResolver.getWriteBufferSizeInBytes(conf));
+
+    conf.set(FlinkOptions.WRITE_BUFFER_MEMORY_TYPE, 
BufferMemoryType.MANAGED.name().toLowerCase());
+    assertTrue(OptionsResolver.isManagedMemoryBufferEnabled(conf));
+  }
+
+  public static class TestPartitioner implements Partitioner<String> {
+    public TestPartitioner(Configuration conf) {
+    }
+
+    @Override
+    public int partition(String key, int numPartitions) {
+      return 0;
+    }
+  }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestFilebasedSchemaProvider.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestFilebasedSchemaProvider.java
new file mode 100644
index 000000000000..ebcbee01d93f
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestFilebasedSchemaProvider.java
@@ -0,0 +1,92 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.schema;
+
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.exception.HoodieIOException;
+
+import org.apache.flink.configuration.Configuration;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Files;
+import java.nio.file.Path;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * Tests for {@link FilebasedSchemaProvider}.
+ */
+class TestFilebasedSchemaProvider {
+
+  private static final String SOURCE_SCHEMA_KEY = 
"hoodie.streamer.schemaprovider.source.schema.file";
+  private static final String TARGET_SCHEMA_KEY = 
"hoodie.streamer.schemaprovider.target.schema.file";
+  private static final String SOURCE_SCHEMA =
+      
"{\"type\":\"record\",\"name\":\"SourceRecord\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"}]}";
+  private static final String TARGET_SCHEMA =
+      
"{\"type\":\"record\",\"name\":\"TargetRecord\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},"
+          + "{\"name\":\"ts\",\"type\":\"long\",\"default\":0}]}";
+
+  @TempDir
+  Path tempDir;
+
+  @Test
+  void testConfigurationConstructorReturnsSourceSchema() throws IOException {
+    Path sourcePath = writeSchema("source.avsc", SOURCE_SCHEMA);
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.SOURCE_AVRO_SCHEMA_PATH, sourcePath.toString());
+
+    FilebasedSchemaProvider provider = new FilebasedSchemaProvider(conf);
+
+    assertEquals("SourceRecord", provider.getSourceSchema().getName());
+    assertEquals("SourceRecord", provider.getTargetSchema().getName());
+  }
+
+  @Test
+  void testTypedPropertiesConstructorReturnsSeparateTargetSchema() throws 
IOException {
+    Path sourcePath = writeSchema("source.avsc", SOURCE_SCHEMA);
+    Path targetPath = writeSchema("target.avsc", TARGET_SCHEMA);
+    TypedProperties props = new TypedProperties();
+    props.setProperty(SOURCE_SCHEMA_KEY, sourcePath.toString());
+    props.setProperty(TARGET_SCHEMA_KEY, targetPath.toString());
+
+    FilebasedSchemaProvider provider = new FilebasedSchemaProvider(props);
+
+    assertEquals("SourceRecord", provider.getSourceSchema().getName());
+    assertEquals("TargetRecord", provider.getTargetSchema().getName());
+  }
+
+  @Test
+  void testReadFailureIsWrapped() {
+    Configuration conf = new Configuration();
+    conf.set(FlinkOptions.SOURCE_AVRO_SCHEMA_PATH, 
tempDir.resolve("missing.avsc").toString());
+
+    assertThrows(HoodieIOException.class, () -> new 
FilebasedSchemaProvider(conf));
+  }
+
+  private Path writeSchema(String fileName, String schema) throws IOException {
+    Path path = tempDir.resolve(fileName);
+    Files.write(path, schema.getBytes(StandardCharsets.UTF_8));
+    return path;
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestSchemaRegistryProvider.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestSchemaRegistryProvider.java
new file mode 100644
index 000000000000..d73fdf1a9bf8
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/schema/TestSchemaRegistryProvider.java
@@ -0,0 +1,127 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.schema;
+
+import org.apache.hudi.common.config.TypedProperties;
+import org.apache.hudi.exception.HoodieIOException;
+
+import org.apache.avro.Schema;
+import org.junit.jupiter.api.Test;
+
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.net.HttpURLConnection;
+import java.nio.charset.StandardCharsets;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+
+/**
+ * Tests for {@link SchemaRegistryProvider}.
+ */
+class TestSchemaRegistryProvider {
+
+  private static final String SOURCE_URL_KEY = 
"hoodie.streamer.schemaprovider.registry.url";
+  private static final String TARGET_URL_KEY = 
"hoodie.streamer.schemaprovider.registry.targetUrl";
+  private static final String SOURCE_SCHEMA =
+      
"{\"type\":\"record\",\"name\":\"SourceRecord\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"}]}";
+  private static final String TARGET_SCHEMA =
+      
"{\"type\":\"record\",\"name\":\"TargetRecord\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},"
+          + "{\"name\":\"ts\",\"type\":\"long\",\"default\":0}]}";
+
+  @Test
+  void testReturnsSourceAndTargetSchemasFromRegistryResponses() {
+    TypedProperties props = new TypedProperties();
+    props.setProperty(SOURCE_URL_KEY, "http://source:secret@localhost/source";);
+    props.setProperty(TARGET_URL_KEY, "http://localhost/target";);
+    StubSchemaRegistryProvider provider = new 
StubSchemaRegistryProvider(props);
+
+    assertEquals("SourceRecord", provider.getSourceSchema().getName());
+    assertEquals("source:secret", provider.authorizationCredentials);
+    assertEquals("TargetRecord", provider.getTargetSchema().getName());
+    assertEquals("TargetRecord", provider.getTargetHoodieSchema().getName());
+  }
+
+  @Test
+  void testAuthorizationHeaderIsBase64Encoded() {
+    TypedProperties props = new TypedProperties();
+    props.setProperty(SOURCE_URL_KEY, "http://localhost/source";);
+    StubSchemaRegistryProvider provider = new 
StubSchemaRegistryProvider(props);
+    HttpURLConnection connection = mock(HttpURLConnection.class);
+
+    provider.setAuthorizationHeader("source:secret", connection);
+
+    verify(connection).setRequestProperty("Authorization", "Basic 
c291cmNlOnNlY3JldA==");
+  }
+
+  @Test
+  void testTargetDefaultsToSourceRegistry() {
+    TypedProperties props = new TypedProperties();
+    props.setProperty(SOURCE_URL_KEY, "http://localhost/source";);
+    StubSchemaRegistryProvider provider = new 
StubSchemaRegistryProvider(props);
+
+    Schema targetSchema = provider.getTargetSchema();
+
+    assertEquals("SourceRecord", targetSchema.getName());
+    assertNotNull(provider.getTargetHoodieSchema());
+  }
+
+  @Test
+  void testRegistryReadFailureIsWrapped() {
+    TypedProperties props = new TypedProperties();
+    props.setProperty(SOURCE_URL_KEY, "http://localhost/failure";);
+    StubSchemaRegistryProvider provider = new 
StubSchemaRegistryProvider(props);
+
+    assertThrows(HoodieIOException.class, provider::getSourceSchema);
+    assertThrows(HoodieIOException.class, provider::getTargetSchema);
+  }
+
+  private static class StubSchemaRegistryProvider extends 
SchemaRegistryProvider {
+    private String authorizationCredentials;
+
+    StubSchemaRegistryProvider(TypedProperties props) {
+      super(props);
+    }
+
+    @Override
+    protected void setAuthorizationHeader(String creds, HttpURLConnection 
connection) {
+      super.setAuthorizationHeader(creds, connection);
+      authorizationCredentials = creds;
+    }
+
+    @Override
+    protected InputStream getStream(HttpURLConnection connection) throws 
IOException {
+      String path = connection.getURL().getPath();
+      if ("/failure".equals(path)) {
+        throw new IOException("schema registry unavailable");
+      }
+      String schema = "/target".equals(path) ? TARGET_SCHEMA : SOURCE_SCHEMA;
+      String response = "{\"schema\":" + quote(schema) + "}";
+      return new 
ByteArrayInputStream(response.getBytes(StandardCharsets.UTF_8));
+    }
+
+    private static String quote(String value) {
+      return "\"" + value.replace("\\", "\\\\").replace("\"", "\\\"") + "\"";
+    }
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/clustering/TestFlinkClusteringConfig.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/clustering/TestFlinkClusteringConfig.java
new file mode 100644
index 000000000000..b4340ceae20b
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/clustering/TestFlinkClusteringConfig.java
@@ -0,0 +1,110 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.sink.clustering;
+
+import org.apache.hudi.configuration.FlinkOptions;
+import org.apache.hudi.util.StreamerUtil;
+import org.apache.hudi.utils.TestConfigurations;
+
+import com.beust.jcommander.JCommander;
+import org.apache.flink.configuration.Configuration;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests for {@link FlinkClusteringConfig}.
+ */
+class TestFlinkClusteringConfig {
+
+  @TempDir
+  Path tempDir;
+
+  @Test
+  void testParseAndDeriveFlinkConfiguration() throws Exception {
+    Configuration tableConf = 
TestConfigurations.getDefaultConf(tempDir.toString());
+    tableConf.set(FlinkOptions.URL_ENCODE_PARTITIONING, true);
+    tableConf.set(FlinkOptions.HIVE_STYLE_PARTITIONING, true);
+    StreamerUtil.initTableIfNotExists(tableConf);
+
+    FlinkClusteringConfig config = new FlinkClusteringConfig();
+    JCommander.newBuilder().addObject(config).build().parse(
+        "--path", tempDir.toString(),
+        "--clustering-delta-commits", "6",
+        "--clustering-tasks", "4",
+        "--clean-retain-commits", "12",
+        "--clean-retain-hours", "36",
+        "--clean-retain-file-versions", "7",
+        "--archive-min-commits", "25",
+        "--archive-max-commits", "40",
+        "--schedule",
+        "--clean-async-enabled",
+        "--plan-partition-filter-mode", "RECENT_DAYS",
+        "--target-file-max-bytes", "1048576",
+        "--small-file-limit", "524288",
+        "--skip-from-latest-partitions", "2",
+        "--sort-columns", "event_ts,order_id",
+        "--sort-memory", "256",
+        "--max-num-groups", "10",
+        "--target-partitions", "5",
+        "--cluster-begin-partition", "2026-01-01",
+        "--cluster-end-partition", "2026-01-31",
+        "--partition-regex-pattern", "2026-01-.*",
+        "--partition-selected", "2026-01-01,2026-01-02",
+        "--hoodie-conf", "hoodie.test.clustering.option=from-cli");
+
+    Configuration conf = FlinkClusteringConfig.toFlinkConfig(config);
+
+    assertEquals(tempDir.toString(), conf.get(FlinkOptions.PATH));
+    assertEquals(6, conf.get(FlinkOptions.CLUSTERING_DELTA_COMMITS));
+    assertEquals(4, conf.get(FlinkOptions.CLUSTERING_TASKS));
+    assertEquals(12, conf.get(FlinkOptions.CLEAN_RETAIN_COMMITS));
+    assertEquals(36, conf.get(FlinkOptions.CLEAN_RETAIN_HOURS));
+    assertEquals(7, conf.get(FlinkOptions.CLEAN_RETAIN_FILE_VERSIONS));
+    assertEquals(25, conf.get(FlinkOptions.ARCHIVE_MIN_COMMITS));
+    assertEquals(40, conf.get(FlinkOptions.ARCHIVE_MAX_COMMITS));
+    assertEquals("RECENT_DAYS", 
conf.get(FlinkOptions.CLUSTERING_PLAN_PARTITION_FILTER_MODE_NAME));
+    assertEquals(1048576L, 
conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_TARGET_FILE_MAX_BYTES));
+    assertEquals(524288L, 
conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_SMALL_FILE_LIMIT));
+    assertEquals(2, 
conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_SKIP_PARTITIONS_FROM_LATEST));
+    assertEquals("event_ts,order_id", 
conf.get(FlinkOptions.CLUSTERING_SORT_COLUMNS));
+    assertEquals(256, conf.get(FlinkOptions.WRITE_SORT_MEMORY));
+    assertEquals(10, conf.get(FlinkOptions.CLUSTERING_MAX_NUM_GROUPS));
+    assertEquals(5, conf.get(FlinkOptions.CLUSTERING_TARGET_PARTITIONS));
+    assertEquals("2026-01-01",
+        
conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_CLUSTER_BEGIN_PARTITION));
+    assertEquals("2026-01-31",
+        conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_CLUSTER_END_PARTITION));
+    assertEquals("2026-01-.*",
+        
conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_PARTITION_REGEX_PATTERN));
+    assertEquals("2026-01-01,2026-01-02",
+        conf.get(FlinkOptions.CLUSTERING_PLAN_STRATEGY_PARTITION_SELECTED));
+    assertTrue(conf.get(FlinkOptions.CLEAN_ASYNC_ENABLED));
+    assertFalse(conf.get(FlinkOptions.CLUSTERING_ASYNC_ENABLED));
+    assertTrue(conf.get(FlinkOptions.CLUSTERING_SCHEDULE_ENABLED));
+    assertTrue(conf.get(FlinkOptions.URL_ENCODE_PARTITIONING));
+    assertTrue(conf.get(FlinkOptions.HIVE_STYLE_PARTITIONING));
+    assertEquals("from-cli", conf.getString("hoodie.test.clustering.option", 
null));
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/TestFlinkCompactionConfig.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/TestFlinkCompactionConfig.java
new file mode 100644
index 000000000000..10ac0f7357dc
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/compact/TestFlinkCompactionConfig.java
@@ -0,0 +1,79 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.sink.compact;
+
+import org.apache.hudi.common.config.HoodieMemoryConfig;
+import org.apache.hudi.common.config.HoodieReaderConfig;
+import org.apache.hudi.configuration.FlinkOptions;
+
+import com.beust.jcommander.JCommander;
+import org.apache.flink.configuration.Configuration;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests for {@link FlinkCompactionConfig}.
+ */
+class TestFlinkCompactionConfig {
+
+  @TempDir
+  Path tempDir;
+
+  @Test
+  void testParseAndDeriveFlinkConfiguration() {
+    FlinkCompactionConfig config = new FlinkCompactionConfig();
+    JCommander.newBuilder().addObject(config).build().parse(
+        "--path", tempDir.toString(),
+        "--compaction-trigger-strategy", FlinkCompactionConfig.NUM_OR_TIME,
+        "--disable-file-group-reader",
+        "--compaction-delta-commits", "8",
+        "--compaction-delta-seconds", "900",
+        "--clean-async-enabled",
+        "--compaction-max-memory", "256",
+        "--compaction-target-io", "2048",
+        "--compaction-tasks", "4",
+        "--schedule",
+        "--spillable_map_path", tempDir.resolve("spill").toString(),
+        "--hoodie-conf", "hoodie.test.compaction.option=from-cli");
+
+    Configuration conf = FlinkCompactionConfig.toFlinkConfig(config);
+
+    assertEquals(tempDir.toString(), conf.get(FlinkOptions.PATH));
+    assertEquals(FlinkCompactionConfig.NUM_OR_TIME, 
conf.get(FlinkOptions.COMPACTION_TRIGGER_STRATEGY));
+    assertEquals(8, conf.get(FlinkOptions.COMPACTION_DELTA_COMMITS));
+    assertEquals(900, conf.get(FlinkOptions.COMPACTION_DELTA_SECONDS));
+    assertEquals(256, conf.get(FlinkOptions.COMPACTION_MAX_MEMORY));
+    assertEquals(256, conf.get(FlinkOptions.WRITE_MERGE_MAX_MEMORY));
+    assertEquals(2048L, conf.get(FlinkOptions.COMPACTION_TARGET_IO));
+    assertEquals(4, conf.get(FlinkOptions.COMPACTION_TASKS));
+    assertTrue(conf.get(FlinkOptions.CLEAN_ASYNC_ENABLED));
+    
assertFalse(conf.get(FlinkOptions.COMPACTION_OPERATION_EXECUTE_ASYNC_ENABLED));
+    assertTrue(conf.get(FlinkOptions.COMPACTION_SCHEDULE_ENABLED));
+    assertEquals(tempDir.resolve("spill").toString(),
+        conf.getString(HoodieMemoryConfig.SPILLABLE_MAP_BASE_PATH.key(), 
null));
+    assertEquals("false", 
conf.getString(HoodieReaderConfig.FILE_GROUP_READER_ENABLED.key(), null));
+    assertEquals("from-cli", conf.getString("hoodie.test.compaction.option", 
null));
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestFlinkStreamerConfig.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestFlinkStreamerConfig.java
new file mode 100644
index 000000000000..446b95558c57
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestFlinkStreamerConfig.java
@@ -0,0 +1,146 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.streamer;
+
+import org.apache.hudi.common.model.WriteOperationType;
+import org.apache.hudi.configuration.FlinkOptions;
+
+import com.beust.jcommander.JCommander;
+import com.beust.jcommander.ParameterException;
+import org.apache.flink.configuration.Configuration;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.io.TempDir;
+
+import java.nio.file.Path;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Tests for {@link FlinkStreamerConfig}.
+ */
+class TestFlinkStreamerConfig {
+
+  @TempDir
+  Path tempDir;
+
+  @Test
+  void testParseAndDeriveFlinkConfiguration() {
+    FlinkStreamerConfig config = parse(
+        "--kafka-topic", "orders",
+        "--kafka-group-id", "flink-writers",
+        "--kafka-bootstrap-servers", "broker:9092",
+        "--target-base-path", tempDir.toString(),
+        "--target-table", "orders_hudi",
+        "--table-type", "merge_on_read",
+        "--op", "INSERT",
+        "--record-key-field", "order_id",
+        "--partition-path-field", "order_date",
+        "--source-ordering-fields", "event_ts,seq_no",
+        "--instant-retry-times", "7",
+        "--instant-retry-interval", "2",
+        "--filter-dupes",
+        "--commit-on-errors",
+        "--metadata-enabled",
+        "--write-rate-limit", "500",
+        "--write-task-num", "6",
+        "--bucket-assign-num", "5",
+        "--index-bootstrap-num", "4",
+        "--source-avro-schema-path", "file:///tmp/source.avsc",
+        "--source-avro-schema", 
"{\"type\":\"record\",\"name\":\"order\",\"fields\":[]}",
+        "--compaction-tasks", "3",
+        "--clustering-tasks", "2",
+        "--hive-sync-enable",
+        "--hive-sync-db", "analytics",
+        "--hive-sync-table", "orders",
+        "--hoodie-conf", 
"hoodie.datasource.write.drop.partition.columns=true");
+
+    Configuration conf = FlinkStreamerConfig.toFlinkConfig(config);
+
+    assertEquals(tempDir.toString(), conf.get(FlinkOptions.PATH));
+    assertEquals("orders_hudi", conf.get(FlinkOptions.TABLE_NAME));
+    assertEquals("MERGE_ON_READ", conf.get(FlinkOptions.TABLE_TYPE));
+    assertEquals(WriteOperationType.INSERT.value(), 
conf.get(FlinkOptions.OPERATION));
+    assertEquals("order_id", conf.get(FlinkOptions.RECORD_KEY_FIELD));
+    assertEquals("order_date", conf.get(FlinkOptions.PARTITION_PATH_FIELD));
+    assertEquals("event_ts,seq_no", conf.get(FlinkOptions.ORDERING_FIELDS));
+    assertEquals(7, conf.get(FlinkOptions.RETRY_TIMES));
+    assertEquals(2_000L, conf.get(FlinkOptions.RETRY_INTERVAL_MS));
+    assertTrue(conf.get(FlinkOptions.PRE_COMBINE));
+    assertTrue(conf.get(FlinkOptions.IGNORE_FAILED));
+    assertTrue(conf.get(FlinkOptions.METADATA_ENABLED));
+    assertEquals(500L, conf.get(FlinkOptions.WRITE_RATE_LIMIT));
+    assertEquals(6, conf.get(FlinkOptions.WRITE_TASKS));
+    assertEquals(5, conf.get(FlinkOptions.BUCKET_ASSIGN_TASKS));
+    assertEquals(4, conf.get(FlinkOptions.INDEX_BOOTSTRAP_TASKS));
+    assertEquals("file:///tmp/source.avsc", 
conf.get(FlinkOptions.SOURCE_AVRO_SCHEMA_PATH));
+    assertEquals(3, conf.get(FlinkOptions.COMPACTION_TASKS));
+    assertEquals(2, conf.get(FlinkOptions.CLUSTERING_TASKS));
+    assertTrue(conf.get(FlinkOptions.HIVE_SYNC_ENABLED));
+    assertEquals("analytics", conf.get(FlinkOptions.HIVE_SYNC_DB));
+    assertEquals("orders", conf.get(FlinkOptions.HIVE_SYNC_TABLE));
+    assertEquals("true",
+        conf.getString("hoodie.datasource.write.drop.partition.columns", 
null));
+  }
+
+  @Test
+  void testCustomKeyGeneratorTakesPrecedence() {
+    FlinkStreamerConfig config = parseRequiredOptions(
+        "--keygen-class", "org.example.CustomKeyGenerator",
+        "--keygen-type", "COMPLEX");
+
+    Configuration conf = FlinkStreamerConfig.toFlinkConfig(config);
+
+    assertEquals("org.example.CustomKeyGenerator", 
conf.get(FlinkOptions.KEYGEN_CLASS_NAME));
+    assertFalse(conf.contains(FlinkOptions.KEYGEN_TYPE));
+  }
+
+  @Test
+  void testRequiredOptionsAndNumericValuesAreValidated() {
+    FlinkStreamerConfig missingRequired = new FlinkStreamerConfig();
+    assertThrows(ParameterException.class,
+        () -> 
JCommander.newBuilder().addObject(missingRequired).build().parse("--kafka-topic",
 "orders"));
+
+    FlinkStreamerConfig invalidRetry = 
parseRequiredOptions("--instant-retry-times", "not-a-number");
+    assertThrows(NumberFormatException.class, () -> 
FlinkStreamerConfig.toFlinkConfig(invalidRetry));
+  }
+
+  private FlinkStreamerConfig parseRequiredOptions(String... additionalArgs) {
+    String[] requiredArgs = {
+        "--kafka-topic", "orders",
+        "--kafka-group-id", "flink-writers",
+        "--kafka-bootstrap-servers", "broker:9092",
+        "--target-base-path", tempDir.toString(),
+        "--target-table", "orders_hudi",
+        "--table-type", "copy_on_write"
+    };
+    String[] args = new String[requiredArgs.length + additionalArgs.length];
+    System.arraycopy(requiredArgs, 0, args, 0, requiredArgs.length);
+    System.arraycopy(additionalArgs, 0, args, requiredArgs.length, 
additionalArgs.length);
+    return parse(args);
+  }
+
+  private static FlinkStreamerConfig parse(String... args) {
+    FlinkStreamerConfig config = new FlinkStreamerConfig();
+    JCommander.newBuilder().addObject(config).build().parse(args);
+    return config;
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestHoodieFlinkStreamer.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestHoodieFlinkStreamer.java
new file mode 100644
index 000000000000..f848faa0d6ee
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/streamer/TestHoodieFlinkStreamer.java
@@ -0,0 +1,224 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hudi.streamer;
+
+import org.apache.hudi.client.model.HoodieFlinkInternalRow;
+import org.apache.hudi.common.util.Option;
+import org.apache.hudi.configuration.OptionsInference;
+import org.apache.hudi.configuration.OptionsResolver;
+import org.apache.hudi.sink.transform.Transformer;
+import org.apache.hudi.sink.utils.Pipelines;
+import org.apache.hudi.util.StreamerUtil;
+import org.apache.hudi.utils.StreamerUtils;
+
+import org.apache.flink.api.common.ExecutionConfig;
+import org.apache.flink.configuration.CheckpointingOptions;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.StateBackendOptions;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.CheckpointConfig;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.table.data.RowData;
+import org.apache.flink.table.types.logical.RowType;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import java.util.concurrent.atomic.AtomicReference;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyList;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.same;
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Tests the argument and pipeline wiring in {@link HoodieFlinkStreamer}.
+ */
+class TestHoodieFlinkStreamer {
+
+  private static final String SOURCE_SCHEMA =
+      
"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"}]}";
+
+  @Test
+  void testAppendPipelineWiringWithTransformer() throws Exception {
+    StreamExecutionEnvironment env = mockEnvironment();
+    DataStream<RowData> source = mock(DataStream.class);
+    DataStream<RowData> transformed = mock(DataStream.class);
+    DataStream<RowData> pipeline = mock(DataStream.class);
+    Transformer transformer = mock(Transformer.class);
+    when(transformer.apply(source)).thenReturn(transformed);
+    AtomicReference<Configuration> envConf = new AtomicReference<>();
+
+    // Static mocks intercept config inference and resolution so this test 
only exercises entry-point wiring.
+    try (MockedStatic<StreamExecutionEnvironment> environments = 
mockStatic(StreamExecutionEnvironment.class);
+         MockedStatic<StreamerUtils> streamerUtils = 
mockStatic(StreamerUtils.class);
+         MockedStatic<StreamerUtil> streamerUtil = 
mockStatic(StreamerUtil.class, CALLS_REAL_METHODS);
+         MockedStatic<OptionsInference> inference = 
mockStatic(OptionsInference.class);
+         MockedStatic<OptionsResolver> resolver = 
mockStatic(OptionsResolver.class);
+         MockedStatic<Pipelines> pipelines = mockStatic(Pipelines.class)) {
+      environments.when(() -> 
StreamExecutionEnvironment.getExecutionEnvironment(any(Configuration.class)))
+          .thenAnswer(invocation -> {
+            envConf.set(invocation.getArgument(0));
+            return env;
+          });
+      streamerUtils.when(() -> StreamerUtils.createKafkaStream(
+          same(env), any(RowType.class), eq("orders"), 
any())).thenReturn(source);
+      streamerUtil.when(() -> 
StreamerUtil.createTransformer(anyList())).thenReturn(Option.of(transformer));
+      resolver.when(() -> 
OptionsResolver.isAppendMode(any())).thenReturn(true);
+      resolver.when(() -> 
OptionsResolver.needsAsyncClustering(any())).thenReturn(true);
+      pipelines.when(() -> Pipelines.append(any(), any(RowType.class), 
same(transformed))).thenReturn(pipeline);
+
+      HoodieFlinkStreamer.main(args(
+          "--table-type", "COPY_ON_WRITE",
+          "--op", "INSERT",
+          "--transformer-class", "org.example.Transformer",
+          "--flink-checkpoint-path", "file:///tmp/checkpoints"));
+
+      assertEquals("HashMapStateBackend", 
envConf.get().get(StateBackendOptions.STATE_BACKEND));
+      assertEquals("filesystem", 
envConf.get().get(CheckpointingOptions.CHECKPOINT_STORAGE));
+      assertEquals("file:///tmp/checkpoints",
+          envConf.get().get(CheckpointingOptions.CHECKPOINTS_DIRECTORY));
+      pipelines.verify(() -> Pipelines.cluster(any(), any(RowType.class), 
same(pipeline)));
+      verify(env).execute("orders_hudi");
+    }
+  }
+
+  @Test
+  void testAppendPipelineFallbackWiring() throws Exception {
+    StreamExecutionEnvironment env = mockEnvironment();
+    DataStream<RowData> source = mock(DataStream.class);
+    DataStream<RowData> pipeline = mock(DataStream.class);
+
+    try (MockedStatic<StreamExecutionEnvironment> environments = 
mockStatic(StreamExecutionEnvironment.class);
+         MockedStatic<StreamerUtils> streamerUtils = 
mockStatic(StreamerUtils.class);
+         MockedStatic<OptionsInference> inference = 
mockStatic(OptionsInference.class);
+         MockedStatic<OptionsResolver> resolver = 
mockStatic(OptionsResolver.class);
+         MockedStatic<Pipelines> pipelines = mockStatic(Pipelines.class)) {
+      environments.when(() -> 
StreamExecutionEnvironment.getExecutionEnvironment(any(Configuration.class)))
+          .thenReturn(env);
+      streamerUtils.when(() -> StreamerUtils.createKafkaStream(
+          same(env), any(RowType.class), eq("orders"), 
any())).thenReturn(source);
+      resolver.when(() -> 
OptionsResolver.isAppendMode(any())).thenReturn(true);
+      resolver.when(() -> 
OptionsResolver.needsAsyncClustering(any())).thenReturn(false);
+      pipelines.when(() -> Pipelines.append(any(), any(RowType.class), 
same(source))).thenReturn(pipeline);
+
+      resolver.when(() -> 
OptionsResolver.isLazyFailedWritesCleaning(any())).thenReturn(true);
+      HoodieFlinkStreamer.main(args("--table-type", "COPY_ON_WRITE", "--op", 
"INSERT"));
+      pipelines.verify(() -> Pipelines.clean(any(), same(pipeline)));
+
+      resolver.when(() -> 
OptionsResolver.isLazyFailedWritesCleaning(any())).thenReturn(false);
+      HoodieFlinkStreamer.main(args("--table-type", "COPY_ON_WRITE", "--op", 
"INSERT"));
+      pipelines.verify(() -> Pipelines.dummySink(same(pipeline)));
+    }
+  }
+
+  @Test
+  void testUpsertPipelineWiringWithCompaction() throws Exception {
+    StreamExecutionEnvironment env = mockEnvironment();
+    DataStream<RowData> source = mock(DataStream.class);
+    DataStream<HoodieFlinkInternalRow> bootstrapped = mock(DataStream.class);
+    DataStream<RowData> pipeline = mock(DataStream.class);
+
+    // Static mocks intercept config inference and resolution so this test 
only exercises entry-point wiring.
+    try (MockedStatic<StreamExecutionEnvironment> environments = 
mockStatic(StreamExecutionEnvironment.class);
+         MockedStatic<StreamerUtils> streamerUtils = 
mockStatic(StreamerUtils.class);
+         MockedStatic<OptionsInference> inference = 
mockStatic(OptionsInference.class);
+         MockedStatic<OptionsResolver> resolver = 
mockStatic(OptionsResolver.class);
+         MockedStatic<Pipelines> pipelines = mockStatic(Pipelines.class)) {
+      environments.when(() -> 
StreamExecutionEnvironment.getExecutionEnvironment(any(Configuration.class)))
+          .thenReturn(env);
+      streamerUtils.when(() -> StreamerUtils.createKafkaStream(
+          same(env), any(RowType.class), eq("orders"), 
any())).thenReturn(source);
+      resolver.when(() -> 
OptionsResolver.isAppendMode(any())).thenReturn(false);
+      resolver.when(() -> 
OptionsResolver.needsAsyncCompaction(any())).thenReturn(true);
+      pipelines.when(() -> Pipelines.bootstrap(any(), any(RowType.class), 
same(source)))
+          .thenReturn(bootstrapped);
+      pipelines.when(() -> Pipelines.hoodieStreamWrite(any(), 
any(RowType.class), same(bootstrapped)))
+          .thenReturn(pipeline);
+
+      HoodieFlinkStreamer.main(args("--table-type", "MERGE_ON_READ", "--op", 
"UPSERT"));
+
+      pipelines.verify(() -> Pipelines.compact(any(), same(pipeline)));
+      verify(env).execute("orders_hudi");
+    }
+  }
+
+  @Test
+  void testUpsertPipelineFallbackWiring() throws Exception {
+    StreamExecutionEnvironment env = mockEnvironment();
+    DataStream<RowData> source = mock(DataStream.class);
+    DataStream<HoodieFlinkInternalRow> bootstrapped = mock(DataStream.class);
+    DataStream<RowData> pipeline = mock(DataStream.class);
+
+    try (MockedStatic<StreamExecutionEnvironment> environments = 
mockStatic(StreamExecutionEnvironment.class);
+         MockedStatic<StreamerUtils> streamerUtils = 
mockStatic(StreamerUtils.class);
+         MockedStatic<OptionsInference> inference = 
mockStatic(OptionsInference.class);
+         MockedStatic<OptionsResolver> resolver = 
mockStatic(OptionsResolver.class);
+         MockedStatic<Pipelines> pipelines = mockStatic(Pipelines.class)) {
+      environments.when(() -> 
StreamExecutionEnvironment.getExecutionEnvironment(any(Configuration.class)))
+          .thenReturn(env);
+      streamerUtils.when(() -> StreamerUtils.createKafkaStream(
+          same(env), any(RowType.class), eq("orders"), 
any())).thenReturn(source);
+      resolver.when(() -> 
OptionsResolver.isAppendMode(any())).thenReturn(false);
+      resolver.when(() -> 
OptionsResolver.needsAsyncCompaction(any())).thenReturn(false);
+      pipelines.when(() -> Pipelines.bootstrap(any(), any(RowType.class), 
same(source)))
+          .thenReturn(bootstrapped);
+      pipelines.when(() -> Pipelines.hoodieStreamWrite(any(), 
any(RowType.class), same(bootstrapped)))
+          .thenReturn(pipeline);
+
+      resolver.when(() -> 
OptionsResolver.needsAsyncCleaning(any())).thenReturn(true);
+      HoodieFlinkStreamer.main(args("--table-type", "MERGE_ON_READ", "--op", 
"UPSERT"));
+      pipelines.verify(() -> Pipelines.clean(any(), same(pipeline)));
+
+      resolver.when(() -> 
OptionsResolver.needsAsyncCleaning(any())).thenReturn(false);
+      HoodieFlinkStreamer.main(args("--table-type", "MERGE_ON_READ", "--op", 
"UPSERT"));
+      pipelines.verify(() -> Pipelines.dummySink(same(pipeline)));
+    }
+  }
+
+  private static StreamExecutionEnvironment mockEnvironment() {
+    StreamExecutionEnvironment env = mock(StreamExecutionEnvironment.class);
+    CheckpointConfig checkpointConfig = mock(CheckpointConfig.class);
+    ExecutionConfig executionConfig = mock(ExecutionConfig.class);
+    when(env.getCheckpointConfig()).thenReturn(checkpointConfig);
+    when(env.getConfig()).thenReturn(executionConfig);
+    when(checkpointConfig.getCheckpointTimeout()).thenReturn(600_000L);
+    return env;
+  }
+
+  private static String[] args(String... additionalArgs) {
+    String[] requiredArgs = {
+        "--kafka-topic", "orders",
+        "--kafka-group-id", "flink-writers",
+        "--kafka-bootstrap-servers", "broker:9092",
+        "--target-base-path", "file:///tmp/orders",
+        "--target-table", "orders_hudi",
+        "--source-avro-schema", SOURCE_SCHEMA
+    };
+    String[] args = new String[requiredArgs.length + additionalArgs.length];
+    System.arraycopy(requiredArgs, 0, args, 0, requiredArgs.length);
+    System.arraycopy(additionalArgs, 0, args, requiredArgs.length, 
additionalArgs.length);
+    return args;
+  }
+}

Reply via email to