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 ab60e0971f53a6378f85098899e67980fb5ffe6b
Author: Shuo Cheng <[email protected]>
AuthorDate: Thu May 28 21:30:11 2026 +0800

    refactor(flink): Refactor RowData writer factory to use HoodieSchema 
(#18873)
    
    (cherry picked from commit 588998fc719392d7fe428986eb4df56e856c8b4d)
---
 .../io/storage/row/HoodieRowDataCreateHandle.java  | 12 +++-
 .../row/HoodieRowDataFileWriterFactory.java        | 64 ++--------------------
 .../row/HoodieRowDataParquetWriteSupport.java      |  6 +-
 .../io/storage/row/RowDataParquetWriteSupport.java |  6 +-
 .../TestHoodieRowDataParquetConfigInjector.java    | 11 +++-
 5 files changed, 31 insertions(+), 68 deletions(-)

diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataCreateHandle.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataCreateHandle.java
index dbccd3975012..e3a8445d6675 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataCreateHandle.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataCreateHandle.java
@@ -35,11 +35,13 @@ import org.apache.hudi.config.HoodieWriteConfig;
 import org.apache.hudi.exception.HoodieException;
 import org.apache.hudi.exception.HoodieIOException;
 import org.apache.hudi.exception.HoodieInsertException;
+import org.apache.hudi.io.storage.HoodieFileWriterFactory;
 import org.apache.hudi.storage.HoodieStorage;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.table.HoodieTable;
 import org.apache.hudi.table.marker.WriteMarkers;
 import org.apache.hudi.table.marker.WriteMarkersFactory;
+import org.apache.hudi.util.HoodieSchemaConverter;
 
 import lombok.extern.slf4j.Slf4j;
 import org.apache.flink.table.data.RowData;
@@ -293,7 +295,13 @@ public class HoodieRowDataCreateHandle implements 
Serializable {
       Path path, HoodieTable hoodieTable, HoodieWriteConfig config, RowType 
rowType, String instantTime)
       throws IOException {
     StoragePath storagePath = new StoragePath(path.toUri());
-    return (HoodieRowDataFileWriter) new 
HoodieRowDataFileWriterFactory(hoodieTable.getStorage())
-        .getFileWriter(instantTime, storagePath, config, rowType, 
hoodieTable.getTaskContextSupplier());
+    return (HoodieRowDataFileWriter) HoodieFileWriterFactory.getFileWriter(
+        instantTime,
+        storagePath,
+        hoodieTable.getStorage(),
+        config,
+        HoodieSchemaConverter.convertToSchema(rowType).getNonNullType(),
+        hoodieTable.getTaskContextSupplier(),
+        HoodieRecord.HoodieRecordType.FLINK);
   }
 }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
index 49671d37f108..9fe3d6d5f327 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataFileWriterFactory.java
@@ -23,7 +23,6 @@ import org.apache.hudi.common.config.HoodieConfig;
 import org.apache.hudi.common.config.HoodieParquetConfig;
 import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.engine.TaskContextSupplier;
-import org.apache.hudi.common.fs.FSUtils;
 import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.table.HoodieTableConfig;
 import org.apache.hudi.common.util.Option;
@@ -38,7 +37,6 @@ import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
 import org.apache.hudi.util.HoodieSchemaConverter;
-import org.apache.hudi.util.RowDataQueryContexts;
 
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.hadoop.conf.Configuration;
@@ -47,10 +45,6 @@ import org.apache.hadoop.fs.FSDataOutputStream;
 import java.io.IOException;
 import java.io.OutputStream;
 
-import static org.apache.hudi.common.model.HoodieFileFormat.HFILE;
-import static org.apache.hudi.common.model.HoodieFileFormat.LANCE;
-import static org.apache.hudi.common.model.HoodieFileFormat.ORC;
-import static org.apache.hudi.common.model.HoodieFileFormat.PARQUET;
 import static org.apache.hudi.common.util.ParquetUtils.getCompressionCodecName;
 
 /**
@@ -62,30 +56,6 @@ public class HoodieRowDataFileWriterFactory extends 
HoodieFileWriterFactory {
     super(storage);
   }
 
-  public HoodieFileWriter getFileWriter(String instantTime, StoragePath 
storagePath, HoodieWriteConfig config, RowType rowType,
-                                     TaskContextSupplier taskContextSupplier) 
throws IOException {
-    final String extension = FSUtils.getFileExtension(storagePath.getName());
-    return getFileWriterByFormat(extension, instantTime, storagePath, config, 
rowType, taskContextSupplier);
-  }
-
-  private  <T, I, K, O> HoodieFileWriter getFileWriterByFormat(
-      String extension, String instantTime, StoragePath path, HoodieConfig 
config, RowType rowType,
-      TaskContextSupplier taskContextSupplier) throws IOException {
-    if (PARQUET.getFileExtension().equals(extension)) {
-      return newParquetFileWriter(instantTime, path, config, rowType, 
taskContextSupplier);
-    }
-    if (HFILE.getFileExtension().equals(extension)) {
-      return newHFileFileWriter(instantTime, path, config, 
HoodieSchemaConverter.convertToSchema(rowType), taskContextSupplier);
-    }
-    if (ORC.getFileExtension().equals(extension)) {
-      return newOrcFileWriter(instantTime, path, config, 
HoodieSchemaConverter.convertToSchema(rowType), taskContextSupplier);
-    }
-    if (LANCE.getFileExtension().equals(extension)) {
-      return newLanceFileWriter(instantTime, path, config, rowType, 
taskContextSupplier);
-    }
-    throw new UnsupportedOperationException(extension + " format not supported 
yet.");
-  }
-
   /**
    * Create a parquet writer on a given OutputStream.
    *
@@ -100,11 +70,9 @@ public class HoodieRowDataFileWriterFactory extends 
HoodieFileWriterFactory {
       OutputStream outputStream,
       HoodieConfig config,
       HoodieSchema schema) throws IOException {
-    //TODO boundary to revisit in follow up to use HoodieSchema directly
-    final RowType rowType = (RowType) 
RowDataQueryContexts.fromSchema(schema).getRowType().getLogicalType();
     HoodieRowDataParquetWriteSupport writeSupport =
         new HoodieRowDataParquetWriteSupport(
-            storage.getConf().unwrapAs(Configuration.class), rowType, null);
+            storage.getConf().unwrapAs(Configuration.class), schema, null);
     return new HoodieRowDataParquetOutputStreamWriter(
         new FSDataOutputStream(outputStream, null), writeSupport, 
getParquetConfig(config, writeSupport));
   }
@@ -127,28 +95,6 @@ public class HoodieRowDataFileWriterFactory extends 
HoodieFileWriterFactory {
       HoodieConfig config,
       HoodieSchema schema,
       TaskContextSupplier taskContextSupplier) throws IOException {
-    //TODO boundary to revisit in follow up to use HoodieSchema directly
-    final RowType rowType = (RowType) 
RowDataQueryContexts.fromSchema(schema).getRowType().getLogicalType();
-    return newParquetFileWriter(instantTime, storagePath, config, rowType, 
taskContextSupplier);
-  }
-
-  /**
-   * Create a parquet RowData writer on a given storage path.
-   *
-   * @param instantTime         instant time to write
-   * @param storagePath         file storage path
-   * @param config              hoodie configuration
-   * @param rowType             rowType of record
-   * @param taskContextSupplier task context supplier
-   *
-   * @return a RowData parquet writer
-   */
-  public HoodieFileWriter newParquetFileWriter(
-      String instantTime,
-      StoragePath storagePath,
-      HoodieConfig config,
-      RowType rowType,
-      TaskContextSupplier taskContextSupplier) throws IOException {
     boolean populateMetaFields = 
config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS);
     boolean withOperation = 
config.getBooleanOrDefault(HoodieWriteConfig.ALLOW_OPERATION_METADATA_FIELD);
 
@@ -160,23 +106,25 @@ public class HoodieRowDataFileWriterFactory extends 
HoodieFileWriterFactory {
     BloomFilter filter = createBloomFilter(hoodieConfig);
     HoodieRowDataParquetWriteSupport writeSupport = 
(HoodieRowDataParquetWriteSupport) ReflectionUtils.loadClass(
         
hoodieConfig.getStringOrDefault(HoodieStorageConfig.HOODIE_PARQUET_FLINK_ROW_DATA_WRITE_SUPPORT_CLASS),
-        new Class<?>[] {Configuration.class, RowType.class, BloomFilter.class},
-        conf, rowType, filter);
+        new Class<?>[] {Configuration.class, HoodieSchema.class, 
BloomFilter.class},
+        conf, schema, filter);
 
     return new HoodieRowDataParquetWriter(storagePath, 
getParquetConfig(hoodieConfig, writeSupport),
         instantTime, taskContextSupplier, populateMetaFields, withOperation);
   }
 
+  @Override
   public HoodieFileWriter newLanceFileWriter(
       String instantTime,
       StoragePath path,
       HoodieConfig config,
-      RowType rowType,
+      HoodieSchema schema,
       TaskContextSupplier taskContextSupplier) {
     boolean populateMetaFields = 
config.getBooleanOrDefault(HoodieTableConfig.POPULATE_META_FIELDS);
     boolean withOperation = 
config.getBooleanOrDefault(HoodieWriteConfig.ALLOW_OPERATION_METADATA_FIELD);
     Option<org.apache.hudi.common.bloom.BloomFilter> bloomFilter = 
enableBloomFilter(populateMetaFields, config)
         ? Option.of(createBloomFilter(config)) : Option.empty();
+    RowType rowType = HoodieSchemaConverter.convertToRowType(schema);
     return new HoodieRowDataLanceWriter(
         path,
         rowType,
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetWriteSupport.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetWriteSupport.java
index 8dfb06872e9c..14ed278b70a2 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetWriteSupport.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/HoodieRowDataParquetWriteSupport.java
@@ -20,10 +20,10 @@ package org.apache.hudi.io.storage.row;
 
 import org.apache.hudi.avro.HoodieBloomFilterWriteSupport;
 import org.apache.hudi.common.bloom.BloomFilter;
+import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.util.Option;
 
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.types.logical.RowType;
 import org.apache.hadoop.conf.Configuration;
 import org.apache.parquet.hadoop.api.WriteSupport;
 
@@ -37,8 +37,8 @@ public class HoodieRowDataParquetWriteSupport extends 
RowDataParquetWriteSupport
 
   private final Option<HoodieBloomFilterWriteSupport<String>> 
bloomFilterWriteSupportOpt;
 
-  public HoodieRowDataParquetWriteSupport(Configuration conf, RowType rowType, 
BloomFilter bloomFilter) {
-    super(rowType, conf);
+  public HoodieRowDataParquetWriteSupport(Configuration conf, HoodieSchema 
schema, BloomFilter bloomFilter) {
+    super(schema, conf);
     this.bloomFilterWriteSupportOpt = Option.ofNullable(bloomFilter)
         .map(HoodieBloomFilterRowDataWriteSupport::new);
   }
diff --git 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
index 7315461db500..01d2806e7633 100644
--- 
a/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
+++ 
b/hudi-client/hudi-flink-client/src/main/java/org/apache/hudi/io/storage/row/RowDataParquetWriteSupport.java
@@ -19,8 +19,10 @@
 package org.apache.hudi.io.storage.row;
 
 import org.apache.hudi.common.config.HoodieStorageConfig;
+import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.io.storage.row.parquet.ParquetRowDataWriter;
 import org.apache.hudi.io.storage.row.parquet.ParquetSchemaConverter;
+import org.apache.hudi.util.HoodieSchemaConverter;
 
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.types.logical.RowType;
@@ -41,9 +43,9 @@ public class RowDataParquetWriteSupport extends 
WriteSupport<RowData> {
   private ParquetRowDataWriter writer;
   protected final Configuration hadoopConf;
 
-  public RowDataParquetWriteSupport(RowType rowType, Configuration config) {
+  public RowDataParquetWriteSupport(HoodieSchema hoodieSchema, Configuration 
config) {
     super();
-    this.rowType = rowType;
+    this.rowType = HoodieSchemaConverter.convertToRowType(hoodieSchema);
     this.hadoopConf = new Configuration(config);
     this.schema = 
ParquetSchemaConverter.convertToParquetMessageType("flink_schema", rowType);
   }
diff --git 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
index f65a2844649b..67bf8e09cb6d 100644
--- 
a/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
+++ 
b/hudi-client/hudi-flink-client/src/test/java/org/apache/hudi/io/storage/row/TestHoodieRowDataParquetConfigInjector.java
@@ -21,6 +21,7 @@ package org.apache.hudi.io.storage.row;
 import org.apache.hudi.common.config.HoodieConfig;
 import org.apache.hudi.common.config.HoodieStorageConfig;
 import org.apache.hudi.common.engine.LocalTaskContextSupplier;
+import org.apache.hudi.common.schema.HoodieSchema;
 import org.apache.hudi.common.testutils.DisableDictionaryInjector;
 import org.apache.hudi.common.testutils.HoodieTestUtils;
 import org.apache.hudi.common.util.collection.Pair;
@@ -31,6 +32,7 @@ import org.apache.hudi.storage.StorageConfiguration;
 import org.apache.hudi.storage.StoragePath;
 import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration;
 import org.apache.hudi.testutils.HoodieFlinkClientTestHarness;
+import org.apache.hudi.util.HoodieSchemaConverter;
 
 import org.apache.flink.table.api.DataTypes;
 import org.apache.flink.table.data.GenericRowData;
@@ -111,6 +113,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
         basePath + "/partition/path/test_dictionary_" + instantTime + 
".parquet");
 
     RowType rowType = getTestRowType();
+    HoodieSchema schema = HoodieSchemaConverter.convertToSchema(rowType);
 
     // Create config with the custom injector
     HoodieConfig config = new HoodieConfig();
@@ -120,7 +123,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
     // Create writer and write some data
     HoodieRowDataFileWriterFactory factory = new 
HoodieRowDataFileWriterFactory(storage);
     HoodieFileWriter writer = factory.newParquetFileWriter(
-        instantTime, parquetPath, config, rowType, new 
LocalTaskContextSupplier());
+        instantTime, parquetPath, config, schema, new 
LocalTaskContextSupplier());
 
     assertTrue(writer instanceof HoodieRowDataParquetWriter);
 
@@ -169,6 +172,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
         basePath + "/partition/path/test_invalid_" + instantTime + ".parquet");
 
     RowType rowType = getTestRowType();
+    HoodieSchema schema = HoodieSchemaConverter.convertToSchema(rowType);
 
     // Create config with an invalid/non-existent injector class
     HoodieConfig config = new HoodieConfig();
@@ -177,7 +181,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
     // Should throw an exception when trying to create the writer
     HoodieRowDataFileWriterFactory factory = new 
HoodieRowDataFileWriterFactory(storage);
     assertThrows(Exception.class, () -> {
-      factory.newParquetFileWriter(instantTime, parquetPath, config, rowType, 
new LocalTaskContextSupplier());
+      factory.newParquetFileWriter(instantTime, parquetPath, config, schema, 
new LocalTaskContextSupplier());
     });
   }
 
@@ -189,6 +193,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
         basePath + "/partition/path/test_no_injector_" + instantTime + 
".parquet");
 
     RowType rowType = getTestRowType();
+    HoodieSchema schema = HoodieSchemaConverter.convertToSchema(rowType);
 
     // Create config WITHOUT injector - should use default settings
     HoodieConfig config = new HoodieConfig();
@@ -197,7 +202,7 @@ public class TestHoodieRowDataParquetConfigInjector extends 
HoodieFlinkClientTes
     // Create writer and write some data
     HoodieRowDataFileWriterFactory factory = new 
HoodieRowDataFileWriterFactory(storage);
     HoodieFileWriter writer = factory.newParquetFileWriter(
-        instantTime, parquetPath, config, rowType, new 
LocalTaskContextSupplier());
+        instantTime, parquetPath, config, schema, new 
LocalTaskContextSupplier());
 
     assertTrue(writer instanceof HoodieRowDataParquetWriter);
 

Reply via email to