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

wombatu-kun 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 727ded72b252 fix: remove the dependency to flink-table-planner (#19131)
727ded72b252 is described below

commit 727ded72b2522efc3e5bc1dc369bf82814d14d75
Author: Danny Chan <[email protected]>
AuthorDate: Thu Jul 2 16:24:25 2026 +0800

    fix: remove the dependency to flink-table-planner (#19131)
---
 hudi-examples/hudi-examples-flink/pom.xml          |  13 +-
 hudi-flink-datasource/hudi-flink/pom.xml           |  13 +-
 .../org/apache/hudi/sink/StreamWriteFunction.java  |   6 +-
 .../AppendWriteFunctionWithBIMBufferSort.java      |   6 +-
 .../AppendWriteFunctionWithContinuousSort.java     |   6 +-
 ...AppendWriteFunctionWithDisruptorBufferSort.java |   6 +-
 .../hudi/sink/bulk/sort/SortOperatorGen.java       | 433 ++++++++++++++++++++-
 .../hudi/sink/clustering/ClusteringOperator.java   |  11 +-
 .../sink/clustering/HoodieFlinkClusteringJob.java  |   5 +-
 .../hudi/sink/utils/FlinkTransformationUtils.java  |  41 ++
 .../java/org/apache/hudi/sink/utils/Pipelines.java |   9 +-
 .../org/apache/hudi/sink/v2/utils/PipelinesV2.java |   4 +-
 .../hudi/sink/bulk/sort/TestSortOperatorGen.java   | 138 +++++++
 .../sink/cluster/ITTestHoodieFlinkClustering.java  |  10 +-
 .../apache/hudi/table/ITTestVectorDataSource.java  |   4 +-
 .../hudi/table/catalog/TestHoodieCatalog.java      |   5 +-
 .../hudi/table/catalog/TestHoodieHiveCatalog.java  |   6 +-
 hudi-flink-datasource/hudi-flink1.18.x/pom.xml     |   6 -
 hudi-flink-datasource/hudi-flink1.19.x/pom.xml     |   6 -
 hudi-flink-datasource/hudi-flink1.20.x/pom.xml     |   6 -
 hudi-flink-datasource/hudi-flink2.0.x/pom.xml      |   6 -
 hudi-flink-datasource/hudi-flink2.1.x/pom.xml      |   6 -
 pom.xml                                            |  13 +-
 23 files changed, 653 insertions(+), 106 deletions(-)

diff --git a/hudi-examples/hudi-examples-flink/pom.xml 
b/hudi-examples/hudi-examples-flink/pom.xml
index 077b449189dc..b23579130704 100644
--- a/hudi-examples/hudi-examples-flink/pom.xml
+++ b/hudi-examples/hudi-examples-flink/pom.xml
@@ -200,12 +200,6 @@
             <version>${flink.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>${flink.table.planner.artifactId}</artifactId>
-            <version>${flink.version}</version>
-            <scope>provided</scope>
-        </dependency>
         <dependency>
             <groupId>org.apache.flink</groupId>
             <artifactId>${flink.statebackend.rocksdb.artifactId}</artifactId>
@@ -377,6 +371,13 @@
             <scope>test</scope>
             <type>test-jar</type>
         </dependency>
+        <!-- Required by TableEnvironmentImpl in TestHoodieFlinkQuickstart to 
discover ExecutorFactory. -->
+        <dependency>
+            <groupId>org.apache.flink</groupId>
+            <artifactId>${flink.table.planner.artifactId}</artifactId>
+            <version>${flink.version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.apache.flink</groupId>
             <artifactId>flink-csv</artifactId>
diff --git a/hudi-flink-datasource/hudi-flink/pom.xml 
b/hudi-flink-datasource/hudi-flink/pom.xml
index 2a32636da7a9..6db5d076aa37 100644
--- a/hudi-flink-datasource/hudi-flink/pom.xml
+++ b/hudi-flink-datasource/hudi-flink/pom.xml
@@ -234,12 +234,6 @@
             <version>${flink.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>${flink.table.planner.artifactId}</artifactId>
-            <version>${flink.version}</version>
-            <scope>provided</scope>
-        </dependency>
         <dependency>
             <groupId>org.apache.flink</groupId>
             <artifactId>${flink.statebackend.rocksdb.artifactId}</artifactId>
@@ -452,6 +446,13 @@
             <scope>test</scope>
             <type>test-jar</type>
         </dependency>
+        <!-- Required by TableEnvironmentImpl in SQL tests to discover 
ExecutorFactory. -->
+        <dependency>
+            <groupId>org.apache.flink</groupId>
+            <artifactId>${flink.table.planner.artifactId}</artifactId>
+            <version>${flink.version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.apache.flink</groupId>
             <artifactId>flink-csv</artifactId>
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
index 073b977aa8e6..51494f6f7615 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/StreamWriteFunction.java
@@ -56,7 +56,6 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
 import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
 import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
 import org.apache.flink.table.runtime.operators.sort.BinaryInMemorySortBuffer;
@@ -220,9 +219,8 @@ public class StreamWriteFunction extends 
AbstractStreamWriteFunction<HoodieFlink
     ValidationUtils.checkArgument(recordKeyFields.length > 0,
         "Record key fields can't be empty for LSM storage layout stream 
write.");
     SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType, 
recordKeyFields);
-    SortCodeGenerator codeGenerator = 
sortOperatorGen.createSortCodeGenerator();
-    this.recordKeyComputer = 
codeGenerator.generateNormalizedKeyComputer("LsmRecordKeySortComputer");
-    this.recordKeyComparator = 
codeGenerator.generateRecordComparator("LsmRecordKeySortComparator");
+    this.recordKeyComputer = 
sortOperatorGen.generateNormalizedKeyComputer("LsmRecordKeySortComputer");
+    this.recordKeyComparator = 
sortOperatorGen.generateRecordComparator("LsmRecordKeySortComparator");
     log.info("LSM storage layout stream write will sort buffered RowData by 
record keys: {}",
         String.join(",", recordKeyFields));
   }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithBIMBufferSort.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithBIMBufferSort.java
index 091f018c3c5f..11017d9e4fc2 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithBIMBufferSort.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithBIMBufferSort.java
@@ -33,7 +33,6 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.runtime.operators.sort.QuickSort;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
 import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
 import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
 import org.apache.flink.table.runtime.operators.sort.BinaryInMemorySortBuffer;
@@ -89,9 +88,8 @@ public class AppendWriteFunctionWithBIMBufferSort<T> extends 
AppendWriteFunction
     // Resolve sort keys (defaults to record key if not specified)
     List<String> sortKeyList = AppendWriteFunctions.resolveSortKeys(config);
     SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType, 
sortKeyList.toArray(new String[0]));
-    SortCodeGenerator codeGenerator = 
sortOperatorGen.createSortCodeGenerator();
-    GeneratedNormalizedKeyComputer keyComputer = 
codeGenerator.generateNormalizedKeyComputer("SortComputer");
-    GeneratedRecordComparator recordComparator = 
codeGenerator.generateRecordComparator("SortComparator");
+    GeneratedNormalizedKeyComputer keyComputer = 
sortOperatorGen.generateNormalizedKeyComputer("SortComputer");
+    GeneratedRecordComparator recordComparator = 
sortOperatorGen.generateRecordComparator("SortComparator");
     this.memorySegmentPools = 
this.memorySegmentPoolFactory.createMemorySegmentPools(config, 2, 
OptionsResolver.getWriteBufferSizeInBytes(config));
 
     this.activeBuffer = BufferUtils.createBuffer(rowType,
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithContinuousSort.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithContinuousSort.java
index 8f6b602c1344..78cddd20037f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithContinuousSort.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithContinuousSort.java
@@ -31,7 +31,6 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.core.memory.MemorySegment;
 import org.apache.flink.core.memory.MemorySegmentFactory;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
 import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
 import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
 import org.apache.flink.table.runtime.generated.NormalizedKeyComputer;
@@ -126,9 +125,8 @@ public class AppendWriteFunctionWithContinuousSort<T> 
extends AppendWriteFunctio
 
     // Create sort code generator for normalized key computation and record 
comparison
     SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType, 
sortKeyList.toArray(new String[0]));
-    SortCodeGenerator codeGenerator = 
sortOperatorGen.createSortCodeGenerator();
-    GeneratedNormalizedKeyComputer generatedKeyComputer = 
codeGenerator.generateNormalizedKeyComputer("ContinuousSortKeyComputer");
-    GeneratedRecordComparator generatedComparator = 
codeGenerator.generateRecordComparator("ContinuousSortComparator");
+    GeneratedNormalizedKeyComputer generatedKeyComputer = 
sortOperatorGen.generateNormalizedKeyComputer("ContinuousSortKeyComputer");
+    GeneratedRecordComparator generatedComparator = 
sortOperatorGen.generateRecordComparator("ContinuousSortComparator");
 
     // Instantiate code-generated components
     ClassLoader classLoader = Thread.currentThread().getContextClassLoader();
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithDisruptorBufferSort.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithDisruptorBufferSort.java
index df99a71b9bbe..347ad81da22d 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithDisruptorBufferSort.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/append/AppendWriteFunctionWithDisruptorBufferSort.java
@@ -35,7 +35,6 @@ import org.apache.flink.configuration.Configuration;
 import org.apache.flink.runtime.operators.sort.QuickSort;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
 import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
 import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
 import org.apache.flink.table.runtime.operators.sort.BinaryInMemorySortBuffer;
@@ -95,9 +94,8 @@ public class AppendWriteFunctionWithDisruptorBufferSort<T> 
extends AppendWriteFu
 
     // Create Flink-native sort components
     SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType, 
sortKeyList.toArray(new String[0]));
-    SortCodeGenerator codeGenerator = 
sortOperatorGen.createSortCodeGenerator();
-    this.keyComputer = 
codeGenerator.generateNormalizedKeyComputer("SortComputer");
-    this.recordComparator = 
codeGenerator.generateRecordComparator("SortComparator");
+    this.keyComputer = 
sortOperatorGen.generateNormalizedKeyComputer("SortComputer");
+    this.recordComparator = 
sortOperatorGen.generateRecordComparator("SortComparator");
     this.memorySegmentPool = 
this.memorySegmentPoolFactory.createMemorySegmentPool(config, 
OptionsResolver.getWriteBufferSizeInBytes(config));
 
     initDisruptorBuffer();
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/sort/SortOperatorGen.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/sort/SortOperatorGen.java
index 068f7540f953..2db08f976456 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/sort/SortOperatorGen.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/bulk/sort/SortOperatorGen.java
@@ -20,40 +20,449 @@ package org.apache.hudi.sink.bulk.sort;
 
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.streaming.api.operators.OneInputStreamOperator;
-import org.apache.flink.table.api.TableConfig;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
-import org.apache.flink.table.planner.plan.nodes.exec.spec.SortSpec;
+import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer;
+import org.apache.flink.table.runtime.generated.GeneratedRecordComparator;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.LocalZonedTimestampType;
+import org.apache.flink.table.types.logical.LogicalType;
 import org.apache.flink.table.types.logical.RowType;
+import org.apache.flink.table.types.logical.TimestampType;
+import org.apache.flink.table.types.logical.ZonedTimestampType;
 
+import java.util.ArrayList;
 import java.util.Arrays;
+import java.util.List;
 
 /**
  * Tools to generate the sort operator.
  */
 public class SortOperatorGen {
+  private static final int MAX_NORMALIZED_KEY_BYTES = 16;
+  private static final int VARIABLE_LENGTH_NORMALIZED_KEY_BYTES = 8;
+
   private final int[] sortIndices;
   private final RowType rowType;
-  private final TableConfig tableConfig = TableConfig.getDefault();
+  private final RowData.FieldGetter[] fieldGetters;
 
   public SortOperatorGen(RowType rowType, String[] sortFields) {
-    this.sortIndices = 
Arrays.stream(sortFields).mapToInt(rowType::getFieldIndex).toArray();
     this.rowType = rowType;
+    this.sortIndices = Arrays.stream(sortFields).mapToInt(field -> {
+      int index = rowType.getFieldIndex(field);
+      if (index < 0) {
+        throw new IllegalArgumentException("Can not find sort field '" + field 
+ "' in row type " + rowType);
+      }
+      return index;
+    }).toArray();
+    this.fieldGetters = Arrays.stream(sortIndices)
+        .mapToObj(index -> RowData.createFieldGetter(rowType.getTypeAt(index), 
index))
+        .toArray(RowData.FieldGetter[]::new);
   }
 
   public OneInputStreamOperator<RowData, RowData> 
createSortOperator(Configuration conf) {
-    SortCodeGenerator codeGen = createSortCodeGenerator();
     return new SortOperator(
-        codeGen.generateNormalizedKeyComputer("SortComputer"),
-        codeGen.generateRecordComparator("SortComparator"),
+        generateNormalizedKeyComputer("SortComputer"),
+        generateRecordComparator("SortComparator"),
         conf);
   }
 
-  public SortCodeGenerator createSortCodeGenerator() {
-    SortSpec.SortSpecBuilder builder = SortSpec.builder();
+  public GeneratedNormalizedKeyComputer generateNormalizedKeyComputer(String 
name) {
+    String className = generatedClassName(name);
+    return new GeneratedNormalizedKeyComputer(className, 
generateNormalizedKeyComputerCode(className));
+  }
+
+  public GeneratedRecordComparator generateRecordComparator(String name) {
+    String className = generatedClassName(name);
+    return new GeneratedRecordComparator(className, 
generateRecordComparatorCode(className), fieldGetters);
+  }
+
+  private String generatedClassName(String name) {
+    String normalizedName = name.replaceAll("[^A-Za-z0-9_$]", "_");
+    if (normalizedName.isEmpty() || 
!Character.isJavaIdentifierStart(normalizedName.charAt(0))) {
+      normalizedName = "_" + normalizedName;
+    }
+    int hash = 31 * Arrays.hashCode(sortIndices) + 
rowType.asSerializableString().hashCode();
+    return normalizedName + "_" + Integer.toUnsignedString(hash);
+  }
+
+  private String generateRecordComparatorCode(String className) {
+    StringBuilder code = new StringBuilder();
+    code.append("public final class ").append(className)
+        .append(" implements 
org.apache.flink.table.runtime.generated.RecordComparator {\n")
+        .append("  private final Object[] references;\n")
+        .append("  public ").append(className).append("(Object[] references) 
{\n")
+        .append("    this.references = references;\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public int compare(org.apache.flink.table.data.RowData 
row1, ")
+        .append("org.apache.flink.table.data.RowData row2) {\n");
+    for (int i = 0; i < sortIndices.length; i++) {
+      int sortIndex = sortIndices[i];
+      code.append("    boolean isNull1_").append(i).append(" = 
row1.isNullAt(").append(sortIndex).append(");\n")
+          .append("    boolean isNull2_").append(i).append(" = 
row2.isNullAt(").append(sortIndex).append(");\n")
+          .append("    if (isNull1_").append(i).append(" || 
isNull2_").append(i).append(") {\n")
+          .append("      if (isNull1_").append(i).append(" && 
isNull2_").append(i).append(") {\n")
+          .append("      } else {\n")
+          .append("        return isNull1_").append(i).append(" ? 1 : -1;\n")
+          .append("      }\n")
+          .append("    } else {\n")
+          .append("      int cmp_").append(i).append(" = 
").append(compareExpression(i, sortIndex)).append(";\n")
+          .append("      if (cmp_").append(i).append(" != 0) {\n")
+          .append("        return cmp_").append(i).append(";\n")
+          .append("      }\n")
+          .append("    }\n");
+    }
+    code.append("    return 0;\n")
+        .append("  }\n")
+        .append("  private int 
compareFallback(org.apache.flink.table.data.RowData row1, ")
+        .append("org.apache.flink.table.data.RowData row2, int referenceIndex) 
{\n")
+        .append("    Object value1 = 
((org.apache.flink.table.data.RowData.FieldGetter) references[referenceIndex])")
+        .append(".getFieldOrNull(row1);\n")
+        .append("    Object value2 = 
((org.apache.flink.table.data.RowData.FieldGetter) references[referenceIndex])")
+        .append(".getFieldOrNull(row2);\n")
+        .append("    return compareValues(value1, value2);\n")
+        .append("  }\n")
+        .append("  private static int compareValues(Object value1, Object 
value2) {\n")
+        .append("    if (value1 == value2) {\n")
+        .append("      return 0;\n")
+        .append("    }\n")
+        .append("    if (value1 == null) {\n")
+        .append("      return 1;\n")
+        .append("    }\n")
+        .append("    if (value2 == null) {\n")
+        .append("      return -1;\n")
+        .append("    }\n")
+        .append("    if (value1 instanceof byte[] && value2 instanceof byte[]) 
{\n")
+        .append("      return compareUnsignedBytes((byte[]) value1, (byte[]) 
value2);\n")
+        .append("    }\n")
+        .append("    if (value1 instanceof java.lang.Comparable && value2 
instanceof java.lang.Comparable) {\n")
+        .append("      return ((java.lang.Comparable) 
value1).compareTo(value2);\n")
+        .append("    }\n")
+        .append("    throw new IllegalArgumentException(\"Unsupported sort 
field value type: \" ")
+        .append("+ value1.getClass().getName());\n")
+        .append("  }\n")
+        .append("  private static int compareUnsignedBytes(byte[] bytes1, 
byte[] bytes2) {\n")
+        .append("    int len = java.lang.Math.min(bytes1.length, 
bytes2.length);\n")
+        .append("    for (int i = 0; i < len; i++) {\n")
+        .append("      int result = java.lang.Byte.toUnsignedInt(bytes1[i]) ")
+        .append("- java.lang.Byte.toUnsignedInt(bytes2[i]);\n")
+        .append("      if (result != 0) {\n")
+        .append("        return result;\n")
+        .append("      }\n")
+        .append("    }\n")
+        .append("    return bytes1.length - bytes2.length;\n")
+        .append("  }\n")
+        .append("}\n");
+    return code.toString();
+  }
+
+  private String compareExpression(int referenceIndex, int sortIndex) {
+    LogicalType logicalType = rowType.getTypeAt(sortIndex);
+    switch (logicalType.getTypeRoot()) {
+      case BOOLEAN:
+        return "java.lang.Boolean.compare(row1.getBoolean(" + sortIndex + "), 
row2.getBoolean(" + sortIndex + "))";
+      case TINYINT:
+        return "java.lang.Byte.compare(row1.getByte(" + sortIndex + "), 
row2.getByte(" + sortIndex + "))";
+      case SMALLINT:
+        return "java.lang.Short.compare(row1.getShort(" + sortIndex + "), 
row2.getShort(" + sortIndex + "))";
+      case INTEGER:
+      case DATE:
+      case TIME_WITHOUT_TIME_ZONE:
+      case INTERVAL_YEAR_MONTH:
+        return "java.lang.Integer.compare(row1.getInt(" + sortIndex + "), 
row2.getInt(" + sortIndex + "))";
+      case BIGINT:
+      case INTERVAL_DAY_TIME:
+        return "java.lang.Long.compare(row1.getLong(" + sortIndex + "), 
row2.getLong(" + sortIndex + "))";
+      case FLOAT:
+        return "java.lang.Float.compare(row1.getFloat(" + sortIndex + "), 
row2.getFloat(" + sortIndex + "))";
+      case DOUBLE:
+        return "java.lang.Double.compare(row1.getDouble(" + sortIndex + "), 
row2.getDouble(" + sortIndex + "))";
+      case CHAR:
+      case VARCHAR:
+        return "row1.getString(" + sortIndex + ").compareTo(row2.getString(" + 
sortIndex + "))";
+      case BINARY:
+      case VARBINARY:
+        return "compareUnsignedBytes(row1.getBinary(" + sortIndex + "), 
row2.getBinary(" + sortIndex + "))";
+      case DECIMAL:
+        DecimalType decimalType = (DecimalType) logicalType;
+        return "row1.getDecimal(" + sortIndex + ", " + 
decimalType.getPrecision() + ", "
+            + decimalType.getScale() + ").compareTo(row2.getDecimal(" + 
sortIndex + ", "
+            + decimalType.getPrecision() + ", " + decimalType.getScale() + 
"))";
+      case TIMESTAMP_WITHOUT_TIME_ZONE:
+        return timestampCompareExpression(sortIndex, ((TimestampType) 
logicalType).getPrecision());
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+        return timestampCompareExpression(sortIndex, 
((LocalZonedTimestampType) logicalType).getPrecision());
+      case TIMESTAMP_WITH_TIME_ZONE:
+        return timestampCompareExpression(sortIndex, ((ZonedTimestampType) 
logicalType).getPrecision());
+      default:
+        return "compareFallback(row1, row2, " + referenceIndex + ")";
+    }
+  }
+
+  private String timestampCompareExpression(int sortIndex, int precision) {
+    return "row1.getTimestamp(" + sortIndex + ", " + precision + 
").compareTo(row2.getTimestamp("
+        + sortIndex + ", " + precision + "))";
+  }
+
+  private String generateNormalizedKeyComputerCode(String className) {
+    List<NormalizedKeyField> normalizedKeyFields = normalizedKeyFields();
+    int numKeyBytes = 
normalizedKeyFields.stream().mapToInt(NormalizedKeyField::totalBytes).sum();
+    boolean fullyDetermines = normalizedKeyFields.size() == sortIndices.length
+        && 
normalizedKeyFields.stream().allMatch(NormalizedKeyField::fullyDetermines);
+
+    StringBuilder code = new StringBuilder();
+    code.append("public final class ").append(className)
+        .append(" implements 
org.apache.flink.table.runtime.generated.NormalizedKeyComputer {\n")
+        .append("  public ").append(className).append("(Object[] references) 
{\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public void putKey(org.apache.flink.table.data.RowData 
rowData, ")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset) {\n");
+    int keyOffset = 0;
+    for (NormalizedKeyField normalizedKeyField : normalizedKeyFields) {
+      int valueOffset = keyOffset + 1;
+      int sortIndex = normalizedKeyField.sortIndex;
+      int valueBytes = normalizedKeyField.valueBytes;
+      code.append("    if (rowData.isNullAt(").append(sortIndex).append(")) 
{\n")
+          .append("      target.put(offset + ").append(keyOffset).append(", 
(byte) 1);\n")
+          .append("      zeroBytes(target, offset + 
").append(valueOffset).append(", ")
+          .append(valueBytes).append(");\n")
+          .append("    } else {\n")
+          .append("      target.put(offset + ").append(keyOffset).append(", 
(byte) 0);\n")
+          .append("      ").append(normalizedKeyExpression(sortIndex, 
valueOffset, valueBytes)).append("\n")
+          .append("    }\n");
+      keyOffset += normalizedKeyField.totalBytes();
+    }
+    code.append("  }\n")
+        .append("  @Override\n")
+        .append("  public int 
compareKey(org.apache.flink.core.memory.MemorySegment memorySegment, int i, ")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset) {\n")
+        .append("    for (int j = 0; j < ").append(numKeyBytes).append("; j++) 
{\n")
+        .append("      int cmp = 
java.lang.Byte.toUnsignedInt(memorySegment.get(i + j))\n")
+        .append("          - java.lang.Byte.toUnsignedInt(target.get(offset + 
j));\n")
+        .append("      if (cmp != 0) {\n")
+        .append("        return cmp;\n")
+        .append("      }\n")
+        .append("    }\n")
+        .append("    return 0;\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public void 
swapKey(org.apache.flink.core.memory.MemorySegment seg1, int index1, ")
+        .append("org.apache.flink.core.memory.MemorySegment seg2, int index2) 
{\n")
+        .append("    for (int j = 0; j < ").append(numKeyBytes).append("; j++) 
{\n")
+        .append("      byte tmp = seg1.get(index1 + j);\n")
+        .append("      seg1.put(index1 + j, seg2.get(index2 + j));\n")
+        .append("      seg2.put(index2 + j, tmp);\n")
+        .append("    }\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public int getNumKeyBytes() {\n")
+        .append("    return ").append(numKeyBytes).append(";\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public boolean isKeyFullyDetermines() {\n")
+        .append("    return ").append(fullyDetermines).append(";\n")
+        .append("  }\n")
+        .append("  @Override\n")
+        .append("  public boolean invertKey() {\n")
+        .append("    return false;\n")
+        .append("  }\n")
+        .append("  private static void 
zeroBytes(org.apache.flink.core.memory.MemorySegment target, ")
+        .append("int offset, int numBytes) {\n")
+        .append("    for (int i = 0; i < numBytes; i++) {\n")
+        .append("      target.put(offset + i, (byte) 0);\n")
+        .append("    }\n")
+        .append("  }\n")
+        .append("  private static void putBytesNormalizedKey(byte[] bytes, ")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset, int numBytes) {\n")
+        .append("    int len = java.lang.Math.min(bytes.length, numBytes);\n")
+        .append("    for (int i = 0; i < len; i++) {\n")
+        .append("      target.put(offset + i, bytes[i]);\n")
+        .append("    }\n")
+        .append("    zeroBytes(target, offset + len, numBytes - len);\n")
+        .append("  }\n")
+        .append("  private static void putFloatNormalizedKey(float value, ")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset, int numBytes) {\n")
+        .append("    int bits = java.lang.Float.floatToIntBits(value);\n")
+        .append("    int normalized = bits >= 0 ? bits ^ 
java.lang.Integer.MIN_VALUE : ~bits;\n")
+        .append("    
org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil")
+        .append(".putUnsignedIntegerNormalizedKey(normalized, target, offset, 
numBytes);\n")
+        .append("  }\n")
+        .append("  private static void putDoubleNormalizedKey(double value, ")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset, int numBytes) {\n")
+        .append("    long bits = java.lang.Double.doubleToLongBits(value);\n")
+        .append("    long normalized = bits >= 0 ? bits ^ 
java.lang.Long.MIN_VALUE : ~bits;\n")
+        .append("    
org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil")
+        .append(".putUnsignedLongNormalizedKey(normalized, target, offset, 
numBytes);\n")
+        .append("  }\n")
+        .append("  private static void 
putTimestampNormalizedKey(org.apache.flink.table.data.TimestampData timestamp, 
")
+        .append("org.apache.flink.core.memory.MemorySegment target, int 
offset, int numBytes) {\n")
+        .append("    
org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil")
+        .append(".putLongNormalizedKey(timestamp.getMillisecond(), target, 
offset, java.lang.Math.min(numBytes, 8));\n")
+        .append("    if (numBytes > 8) {\n")
+        .append("      
org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil")
+        .append(".putIntNormalizedKey(timestamp.getNanoOfMillisecond(), 
target, offset + 8, numBytes - 8);\n")
+        .append("    }\n")
+        .append("  }\n")
+        .append("}\n");
+    return code.toString();
+  }
+
+  private List<NormalizedKeyField> normalizedKeyFields() {
+    List<NormalizedKeyField> normalizedKeyFields = new ArrayList<>();
+    int remainingBytes = MAX_NORMALIZED_KEY_BYTES;
     for (int sortIndex : sortIndices) {
-      builder.addField(sortIndex, true, true);
+      LogicalType logicalType = rowType.getTypeAt(sortIndex);
+      int maxValueBytes = maxNormalizedKeyValueBytes(logicalType);
+      if (maxValueBytes <= 0 || remainingBytes <= 1) {
+        break;
+      }
+
+      int valueBytes = Math.min(maxValueBytes, remainingBytes - 1);
+      boolean fullyDetermines = isFixedLengthNormalizedKey(logicalType) && 
valueBytes == maxValueBytes;
+      normalizedKeyFields.add(new NormalizedKeyField(sortIndex, valueBytes, 
fullyDetermines));
+      remainingBytes -= valueBytes + 1;
+      if (!fullyDetermines) {
+        break;
+      }
+    }
+    return normalizedKeyFields;
+  }
+
+  private int maxNormalizedKeyValueBytes(LogicalType logicalType) {
+    switch (logicalType.getTypeRoot()) {
+      case BOOLEAN:
+      case TINYINT:
+        return 1;
+      case SMALLINT:
+        return 2;
+      case INTEGER:
+      case DATE:
+      case TIME_WITHOUT_TIME_ZONE:
+      case INTERVAL_YEAR_MONTH:
+      case FLOAT:
+        return 4;
+      case BIGINT:
+      case INTERVAL_DAY_TIME:
+      case DOUBLE:
+        return 8;
+      case CHAR:
+      case VARCHAR:
+      case BINARY:
+      case VARBINARY:
+        return VARIABLE_LENGTH_NORMALIZED_KEY_BYTES;
+      case DECIMAL:
+        return 
org.apache.flink.table.data.DecimalData.isCompact(((DecimalType) 
logicalType).getPrecision()) ? 8 : 0;
+      case TIMESTAMP_WITHOUT_TIME_ZONE:
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+      case TIMESTAMP_WITH_TIME_ZONE:
+        return 12;
+      default:
+        return 0;
+    }
+  }
+
+  private boolean isFixedLengthNormalizedKey(LogicalType logicalType) {
+    switch (logicalType.getTypeRoot()) {
+      case BOOLEAN:
+      case TINYINT:
+      case SMALLINT:
+      case INTEGER:
+      case DATE:
+      case TIME_WITHOUT_TIME_ZONE:
+      case INTERVAL_YEAR_MONTH:
+      case FLOAT:
+      case BIGINT:
+      case INTERVAL_DAY_TIME:
+      case DOUBLE:
+      case DECIMAL:
+      case TIMESTAMP_WITHOUT_TIME_ZONE:
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+      case TIMESTAMP_WITH_TIME_ZONE:
+        return true;
+      default:
+        return false;
+    }
+  }
+
+  private String normalizedKeyExpression(int sortIndex, int valueOffset, int 
valueBytes) {
+    LogicalType logicalType = rowType.getTypeAt(sortIndex);
+    String offset = "offset + " + valueOffset;
+    switch (logicalType.getTypeRoot()) {
+      case BOOLEAN:
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putBooleanNormalizedKey("
+            + "rowData.getBoolean(" + sortIndex + "), target, " + offset + ", 
" + valueBytes + ");";
+      case TINYINT:
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putByteNormalizedKey("
+            + "rowData.getByte(" + sortIndex + "), target, " + offset + ", " + 
valueBytes + ");";
+      case SMALLINT:
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putShortNormalizedKey("
+            + "rowData.getShort(" + sortIndex + "), target, " + offset + ", " 
+ valueBytes + ");";
+      case INTEGER:
+      case DATE:
+      case TIME_WITHOUT_TIME_ZONE:
+      case INTERVAL_YEAR_MONTH:
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putIntNormalizedKey("
+            + "rowData.getInt(" + sortIndex + "), target, " + offset + ", " + 
valueBytes + ");";
+      case BIGINT:
+      case INTERVAL_DAY_TIME:
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putLongNormalizedKey("
+            + "rowData.getLong(" + sortIndex + "), target, " + offset + ", " + 
valueBytes + ");";
+      case FLOAT:
+        return "putFloatNormalizedKey(rowData.getFloat(" + sortIndex + "), 
target, " + offset + ", "
+            + valueBytes + ");";
+      case DOUBLE:
+        return "putDoubleNormalizedKey(rowData.getDouble(" + sortIndex + "), 
target, " + offset + ", "
+            + valueBytes + ");";
+      case CHAR:
+      case VARCHAR:
+        return "putBytesNormalizedKey(rowData.getString(" + sortIndex + 
").toBytes(), target, " + offset + ", "
+            + valueBytes + ");";
+      case BINARY:
+      case VARBINARY:
+        return "putBytesNormalizedKey(rowData.getBinary(" + sortIndex + "), 
target, " + offset + ", "
+            + valueBytes + ");";
+      case DECIMAL:
+        DecimalType decimalType = (DecimalType) logicalType;
+        return 
"org.apache.flink.api.common.typeutils.base.NormalizedKeyUtil.putLongNormalizedKey("
+            + "rowData.getDecimal(" + sortIndex + ", " + 
decimalType.getPrecision() + ", "
+            + decimalType.getScale() + ").toUnscaledLong(), target, " + offset 
+ ", " + valueBytes + ");";
+      case TIMESTAMP_WITHOUT_TIME_ZONE:
+        return timestampNormalizedKeyExpression(sortIndex, ((TimestampType) 
logicalType).getPrecision(), offset,
+            valueBytes);
+      case TIMESTAMP_WITH_LOCAL_TIME_ZONE:
+        return timestampNormalizedKeyExpression(sortIndex, 
((LocalZonedTimestampType) logicalType).getPrecision(),
+            offset, valueBytes);
+      case TIMESTAMP_WITH_TIME_ZONE:
+        return timestampNormalizedKeyExpression(sortIndex, 
((ZonedTimestampType) logicalType).getPrecision(), offset,
+            valueBytes);
+      default:
+        throw new IllegalArgumentException("Unsupported normalized key field 
type: " + logicalType);
+    }
+  }
+
+  private String timestampNormalizedKeyExpression(int sortIndex, int 
precision, String offset, int valueBytes) {
+    return "putTimestampNormalizedKey(rowData.getTimestamp(" + sortIndex + ", 
" + precision + "), target, "
+        + offset + ", " + valueBytes + ");";
+  }
+
+  private static class NormalizedKeyField {
+    private final int sortIndex;
+    private final int valueBytes;
+    private final boolean fullyDetermines;
+
+    private NormalizedKeyField(int sortIndex, int valueBytes, boolean 
fullyDetermines) {
+      this.sortIndex = sortIndex;
+      this.valueBytes = valueBytes;
+      this.fullyDetermines = fullyDetermines;
+    }
+
+    private int totalBytes() {
+      return valueBytes + 1;
+    }
+
+    private boolean fullyDetermines() {
+      return fullyDetermines;
     }
-    return new SortCodeGenerator(tableConfig, 
Thread.currentThread().getContextClassLoader(), rowType, builder.build());
   }
 }
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
index 36349acbe031..442aa1be68bf 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/ClusteringOperator.java
@@ -74,7 +74,6 @@ import 
org.apache.flink.streaming.runtime.tasks.ProcessingTimeService;
 import org.apache.flink.streaming.runtime.tasks.StreamTask;
 import org.apache.flink.table.data.RowData;
 import org.apache.flink.table.data.binary.BinaryRowData;
-import org.apache.flink.table.planner.codegen.sort.SortCodeGenerator;
 import org.apache.flink.table.runtime.generated.NormalizedKeyComputer;
 import org.apache.flink.table.runtime.generated.RecordComparator;
 import org.apache.flink.table.runtime.operators.TableStreamOperator;
@@ -325,8 +324,9 @@ public class ClusteringOperator extends 
TableStreamOperator<ClusteringCommitEven
 
   private BinaryExternalSorter initSorter() {
     ClassLoader cl = getContainingTask().getUserCodeClassLoader();
-    NormalizedKeyComputer computer = 
createSortCodeGenerator().generateNormalizedKeyComputer("SortComputer").newInstance(cl);
-    RecordComparator comparator = 
createSortCodeGenerator().generateRecordComparator("SortComparator").newInstance(cl);
+    SortOperatorGen sortOperatorGen = createSortOperatorGen();
+    NormalizedKeyComputer computer = 
sortOperatorGen.generateNormalizedKeyComputer("SortComputer").newInstance(cl);
+    RecordComparator comparator = 
sortOperatorGen.generateRecordComparator("SortComparator").newInstance(cl);
 
     MemoryManager memManager = 
getContainingTask().getEnvironment().getMemoryManager();
     BinaryExternalSorter sorter = Utils.getBinaryExternalSorter(
@@ -348,10 +348,9 @@ public class ClusteringOperator extends 
TableStreamOperator<ClusteringCommitEven
     return sorter;
   }
 
-  private SortCodeGenerator createSortCodeGenerator() {
-    SortOperatorGen sortOperatorGen = new SortOperatorGen(rowType,
+  private SortOperatorGen createSortOperatorGen() {
+    return new SortOperatorGen(rowType,
         conf.get(FlinkOptions.CLUSTERING_SORT_COLUMNS).split(","));
-    return sortOperatorGen.createSortCodeGenerator();
   }
 
   private String getFileIds(List<ClusteringOperation> clusteringOperations) {
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/HoodieFlinkClusteringJob.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/HoodieFlinkClusteringJob.java
index 9a2b07a7b81b..d2f6419515ac 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/HoodieFlinkClusteringJob.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/clustering/HoodieFlinkClusteringJob.java
@@ -49,7 +49,6 @@ import 
org.apache.flink.client.deployment.application.ApplicationExecutionExcept
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
-import org.apache.flink.table.planner.plan.nodes.exec.utils.ExecNodeUtil;
 import org.apache.flink.table.types.DataType;
 import org.apache.flink.table.types.logical.RowType;
 import org.slf4j.Logger;
@@ -60,6 +59,8 @@ import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 
+import static 
org.apache.hudi.sink.utils.FlinkTransformationUtils.setManagedMemoryWeight;
+
 /**
  * Flink hudi clustering program that can be executed manually.
  */
@@ -398,7 +399,7 @@ public class HoodieFlinkClusteringJob {
           .setParallelism(clusteringParallelism);
 
       if (OptionsResolver.sortClusteringEnabled(conf)) {
-        ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+        setManagedMemoryWeight(dataStream.getTransformation(),
             conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
       }
 
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/FlinkTransformationUtils.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/FlinkTransformationUtils.java
new file mode 100644
index 000000000000..c19a9220f71b
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/FlinkTransformationUtils.java
@@ -0,0 +1,41 @@
+/*
+ * 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.utils;
+
+import org.apache.flink.api.dag.Transformation;
+import org.apache.flink.core.memory.ManagedMemoryUseCase;
+
+/**
+ * Utilities for Flink transformations.
+ */
+public final class FlinkTransformationUtils {
+  private FlinkTransformationUtils() {
+  }
+
+  public static <T> void setManagedMemoryWeight(Transformation<T> 
transformation, long memoryBytes) {
+    if (memoryBytes <= 0) {
+      return;
+    }
+    int weight = Math.max(1, (int) (memoryBytes >> 20)); // bytes to MiB
+    
transformation.declareManagedMemoryUseCaseAtOperatorScope(ManagedMemoryUseCase.OPERATOR,
 weight)
+        .ifPresent(previousWeight -> {
+          throw new IllegalStateException("Managed memory weight has been set, 
this should not happen.");
+        });
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
index 7caacf49cf64..59612323b7db 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/utils/Pipelines.java
@@ -80,7 +80,6 @@ import 
org.apache.flink.streaming.api.operators.KeyedProcessOperator;
 import org.apache.flink.streaming.api.operators.ProcessOperator;
 import org.apache.flink.streaming.api.transformations.OneInputTransformation;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.planner.plan.nodes.exec.utils.ExecNodeUtil;
 import org.apache.flink.table.runtime.typeutils.InternalTypeInfo;
 import org.apache.flink.table.types.logical.RowType;
 
@@ -160,7 +159,7 @@ public class Pipelines {
         SortOperatorGen sortOperatorGen = 
BucketBulkInsertWriterHelper.getFileIdSorterGen(rowTypeWithFileId);
         dataStream = dataStream.transform("file_sorter", typeInfo, 
sortOperatorGen.createSortOperator(conf))
             .setParallelism(PARALLELISM_VALUE);
-        ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+        
FlinkTransformationUtils.setManagedMemoryWeight(dataStream.getTransformation(),
             conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
       }
     } else if (!FlinkOptions.isDefaultValueDefined(conf, 
FlinkOptions.PARTITION_PATH_FIELD)) {
@@ -190,7 +189,7 @@ public class Pipelines {
             .transform(isNeededSortInput ? "sorter:(partition_key, 
record_key)" : "sorter:(partition_key)",
                 InternalTypeInfo.of(rowType), 
sortOperatorGen.createSortOperator(conf))
             .setParallelism(PARALLELISM_VALUE);
-        ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+        
FlinkTransformationUtils.setManagedMemoryWeight(dataStream.getTransformation(),
             conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
       }
     }
@@ -568,7 +567,7 @@ public class Pipelines {
             new ClusteringOperator(conf, rowType))
         .setParallelism(conf.get(FlinkOptions.CLUSTERING_TASKS));
     if (OptionsResolver.sortClusteringEnabled(conf)) {
-      ExecNodeUtil.setManagedMemoryWeight(clusteringStream.getTransformation(),
+      
FlinkTransformationUtils.setManagedMemoryWeight(clusteringStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
     }
     DataStreamSink<ClusteringCommitEvent> clusteringCommitEventDataStream = 
clusteringStream.addSink(new ClusteringCommitSink(conf))
@@ -611,7 +610,7 @@ public class Pipelines {
 
   public static void declareManagedMemoryIfNecessary(Configuration conf, 
DataStream<?> dataStream, Supplier<Long> bufferSizeSupplier) {
     if (OptionsResolver.isManagedMemoryBufferEnabled(conf)) {
-      ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(), 
bufferSizeSupplier.get());
+      
FlinkTransformationUtils.setManagedMemoryWeight(dataStream.getTransformation(), 
bufferSizeSupplier.get());
     }
   }
 
diff --git 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/utils/PipelinesV2.java
 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/utils/PipelinesV2.java
index 0bae64ccbbc3..5cf19bfd958a 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/utils/PipelinesV2.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/main/java/org/apache/hudi/sink/v2/utils/PipelinesV2.java
@@ -44,9 +44,9 @@ import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.datastream.DataStreamSink;
 import org.apache.flink.streaming.api.operators.ProcessOperator;
 import org.apache.flink.table.data.RowData;
-import org.apache.flink.table.planner.plan.nodes.exec.utils.ExecNodeUtil;
 import org.apache.flink.table.types.logical.RowType;
 
+import static 
org.apache.hudi.sink.utils.FlinkTransformationUtils.setManagedMemoryWeight;
 import static org.apache.hudi.sink.utils.Pipelines.opUID;
 
 /**
@@ -227,7 +227,7 @@ public class PipelinesV2 {
             new ClusteringOperator(conf, rowType))
         .setParallelism(conf.get(FlinkOptions.CLUSTERING_TASKS));
     if (OptionsResolver.sortClusteringEnabled(conf)) {
-      ExecNodeUtil.setManagedMemoryWeight(clusteringStream.getTransformation(),
+      setManagedMemoryWeight(clusteringStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
     }
     return clusteringStream.transform(
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/sort/TestSortOperatorGen.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/sort/TestSortOperatorGen.java
new file mode 100644
index 000000000000..8a6c7169e73b
--- /dev/null
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/bulk/sort/TestSortOperatorGen.java
@@ -0,0 +1,138 @@
+/*
+ * 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.bulk.sort;
+
+import org.apache.flink.core.memory.MemorySegment;
+import org.apache.flink.core.memory.MemorySegmentFactory;
+import org.apache.flink.table.data.DecimalData;
+import org.apache.flink.table.data.GenericRowData;
+import org.apache.flink.table.data.StringData;
+import org.apache.flink.table.data.TimestampData;
+import org.apache.flink.table.runtime.generated.NormalizedKeyComputer;
+import org.apache.flink.table.runtime.generated.RecordComparator;
+import org.apache.flink.table.types.logical.DecimalType;
+import org.apache.flink.table.types.logical.IntType;
+import org.apache.flink.table.types.logical.LogicalType;
+import org.apache.flink.table.types.logical.RowType;
+import org.apache.flink.table.types.logical.TimestampType;
+import org.apache.flink.table.types.logical.VarBinaryType;
+import org.apache.flink.table.types.logical.VarCharType;
+import org.junit.jupiter.api.Test;
+
+import java.math.BigDecimal;
+
+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 SortOperatorGen}.
+ */
+class TestSortOperatorGen {
+
+  @Test
+  void testGeneratedRecordComparator() {
+    RowType rowType = RowType.of(
+        new LogicalType[] {
+            new IntType(),
+            new VarCharType(),
+            new VarBinaryType(),
+            new DecimalType(10, 2),
+            new TimestampType(3)
+        },
+        new String[] {"id", "name", "bytes", "amount", "ts"});
+
+    SortOperatorGen sortOperatorGen =
+        new SortOperatorGen(rowType, new String[] {"id", "name", "bytes", 
"amount", "ts"});
+    RecordComparator comparator = 
sortOperatorGen.generateRecordComparator("TestSortComparator")
+        .newInstance(Thread.currentThread().getContextClassLoader());
+
+    GenericRowData row1 = row(1, "a", new byte[] {1, 2}, "1.00", 1000);
+    GenericRowData row2 = row(1, "a", new byte[] {1, 3}, "1.00", 1000);
+    GenericRowData row3 = row(1, "a", new byte[] {1, 3}, "2.00", 1000);
+    GenericRowData row4 = row(1, "a", new byte[] {1, 3}, "2.00", 2000);
+    GenericRowData nullRow = row(null, "a", new byte[] {1, 3}, "2.00", 2000);
+
+    assertTrue(comparator.compare(row1, row2) < 0);
+    assertTrue(comparator.compare(row2, row3) < 0);
+    assertTrue(comparator.compare(row3, row4) < 0);
+    assertTrue(comparator.compare(nullRow, row1) > 0);
+    assertTrue(comparator.compare(row1, nullRow) < 0);
+    assertEquals(0, comparator.compare(row4, row4));
+  }
+
+  @Test
+  void testGeneratedNormalizedKeyComputer() {
+    RowType rowType = RowType.of(new LogicalType[] {new IntType()}, new 
String[] {"id"});
+    NormalizedKeyComputer computer = new SortOperatorGen(rowType, new String[] 
{"id"})
+        .generateNormalizedKeyComputer("TestSortComputer")
+        .newInstance(Thread.currentThread().getContextClassLoader());
+    MemorySegment segment1 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+    MemorySegment segment2 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+
+    computer.putKey(GenericRowData.of(1), segment1, 0);
+    computer.putKey(GenericRowData.of(2), segment2, 0);
+
+    assertTrue(computer.getNumKeyBytes() > 1);
+    assertTrue(computer.isKeyFullyDetermines());
+    assertTrue(computer.compareKey(segment1, 0, segment2, 0) < 0);
+  }
+
+  @Test
+  void testGeneratedNormalizedKeyComputerWithNullsLast() {
+    RowType rowType = RowType.of(new LogicalType[] {new IntType()}, new 
String[] {"id"});
+    NormalizedKeyComputer computer = new SortOperatorGen(rowType, new String[] 
{"id"})
+        .generateNormalizedKeyComputer("TestSortComputer")
+        .newInstance(Thread.currentThread().getContextClassLoader());
+    MemorySegment segment1 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+    MemorySegment segment2 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+
+    computer.putKey(GenericRowData.of(1), segment1, 0);
+    computer.putKey(GenericRowData.of((Object) null), segment2, 0);
+
+    assertTrue(computer.compareKey(segment1, 0, segment2, 0) < 0);
+  }
+
+  @Test
+  void testGeneratedNormalizedKeyComputerStopsAfterVariableLengthPrefix() {
+    RowType rowType = RowType.of(
+        new LogicalType[] {new VarCharType(), new IntType()},
+        new String[] {"name", "id"});
+    NormalizedKeyComputer computer = new SortOperatorGen(rowType, new String[] 
{"name", "id"})
+        .generateNormalizedKeyComputer("TestSortComputer")
+        .newInstance(Thread.currentThread().getContextClassLoader());
+    MemorySegment segment1 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+    MemorySegment segment2 = MemorySegmentFactory.wrap(new 
byte[computer.getNumKeyBytes()]);
+
+    computer.putKey(GenericRowData.of(StringData.fromString("abcdefghx"), 2), 
segment1, 0);
+    computer.putKey(GenericRowData.of(StringData.fromString("abcdefghy"), 1), 
segment2, 0);
+
+    assertFalse(computer.isKeyFullyDetermines());
+    assertEquals(0, computer.compareKey(segment1, 0, segment2, 0));
+  }
+
+  private static GenericRowData row(Integer id, String name, byte[] bytes, 
String amount, long timestampMillis) {
+    return GenericRowData.of(
+        id,
+        StringData.fromString(name),
+        bytes,
+        DecimalData.fromBigDecimal(new BigDecimal(amount), 10, 2),
+        TimestampData.fromEpochMillis(timestampMillis));
+  }
+}
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/cluster/ITTestHoodieFlinkClustering.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/cluster/ITTestHoodieFlinkClustering.java
index 4fbf90e2e421..c7816afc38ae 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/cluster/ITTestHoodieFlinkClustering.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/sink/cluster/ITTestHoodieFlinkClustering.java
@@ -65,7 +65,6 @@ import 
org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
 import org.apache.flink.table.api.config.ExecutionConfigOptions;
 import org.apache.flink.table.api.config.TableConfigOptions;
 import org.apache.flink.table.api.internal.TableEnvironmentImpl;
-import org.apache.flink.table.planner.plan.nodes.exec.utils.ExecNodeUtil;
 import org.apache.flink.table.types.DataType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.types.Row;
@@ -84,6 +83,7 @@ import java.util.concurrent.TimeUnit;
 import java.util.stream.Collectors;
 
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR;
+import static 
org.apache.hudi.sink.utils.FlinkTransformationUtils.setManagedMemoryWeight;
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
@@ -193,7 +193,7 @@ public class ITTestHoodieFlinkClustering {
               new ClusteringOperator(conf, rowType))
           .setParallelism(clusteringPlan.getInputGroups().size());
 
-      ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+      setManagedMemoryWeight(dataStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
 
       dataStream
@@ -398,7 +398,7 @@ public class ITTestHoodieFlinkClustering {
                   new ClusteringOperator(conf, rowType))
               .setParallelism(clusteringPlan.getInputGroups().size());
 
-      ExecNodeUtil.setManagedMemoryWeight(
+      setManagedMemoryWeight(
           dataStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
 
@@ -661,7 +661,7 @@ public class ITTestHoodieFlinkClustering {
               new ClusteringOperator(conf, rowType))
           .setParallelism(clusteringPlan.getInputGroups().size());
 
-      ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+      setManagedMemoryWeight(dataStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
 
       dataStream
@@ -765,7 +765,7 @@ public class ITTestHoodieFlinkClustering {
               new ClusteringOperator(conf, rowType))
           .setParallelism(clusteringPlan.getInputGroups().size());
 
-      ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+      setManagedMemoryWeight(dataStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
 
       dataStream
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestVectorDataSource.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestVectorDataSource.java
index 2a19a48cf024..77416c59184f 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestVectorDataSource.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/ITTestVectorDataSource.java
@@ -58,7 +58,6 @@ import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 import org.apache.flink.table.api.TableEnvironment;
 import org.apache.flink.table.api.TableResult;
-import org.apache.flink.table.planner.plan.nodes.exec.utils.ExecNodeUtil;
 import org.apache.flink.table.types.DataType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.types.Row;
@@ -86,6 +85,7 @@ import java.util.concurrent.ExecutionException;
 import java.util.stream.Stream;
 
 import static 
org.apache.hudi.common.testutils.HoodieTestUtils.INSTANT_GENERATOR;
+import static 
org.apache.hudi.sink.utils.FlinkTransformationUtils.setManagedMemoryWeight;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
@@ -561,7 +561,7 @@ public class ITTestVectorDataSource {
               new ClusteringOperator(conf, rowType))
           .setParallelism(clusteringPlan.getInputGroups().size());
 
-      ExecNodeUtil.setManagedMemoryWeight(dataStream.getTransformation(),
+      setManagedMemoryWeight(dataStream.getTransformation(),
           conf.get(FlinkOptions.WRITE_SORT_MEMORY) * 1024L * 1024L);
 
       dataStream
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieCatalog.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieCatalog.java
index 379a769ec815..2edc2f1255a5 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieCatalog.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieCatalog.java
@@ -46,7 +46,6 @@ import org.apache.hudi.utils.CatalogUtils;
 import org.apache.hudi.utils.TestConfigurations;
 import org.apache.hudi.utils.TestData;
 
-import org.apache.flink.calcite.shaded.com.google.common.collect.Lists;
 import org.apache.flink.configuration.Configuration;
 import org.apache.flink.core.fs.Path;
 import org.apache.flink.table.api.DataTypes;
@@ -283,7 +282,7 @@ public class TestHoodieCatalog extends 
BaseTestHoodieCatalog {
     final ResolvedCatalogTable singleKeyMultiplePartitionTable = new 
ResolvedCatalogTable(
         CatalogUtils.createCatalogTable(
             
Schema.newBuilder().fromResolvedSchema(CREATE_TABLE_SCHEMA).build(),
-            Lists.newArrayList("par1", "par2"),
+            Arrays.asList("par1", "par2"),
             EXPECTED_OPTIONS,
             "test"),
         CREATE_TABLE_SCHEMA
@@ -301,7 +300,7 @@ public class TestHoodieCatalog extends 
BaseTestHoodieCatalog {
     final ResolvedCatalogTable multipleKeySinglePartitionTable = new 
ResolvedCatalogTable(
         CatalogUtils.createCatalogTable(
             
Schema.newBuilder().fromResolvedSchema(CREATE_MULTI_KEY_TABLE_SCHEMA).build(),
-            Lists.newArrayList("par1"),
+            Collections.singletonList("par1"),
             EXPECTED_OPTIONS,
             "test"),
         CREATE_TABLE_SCHEMA
diff --git 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieHiveCatalog.java
 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieHiveCatalog.java
index 74e7a1b8ca98..c03c0ad926c0 100644
--- 
a/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieHiveCatalog.java
+++ 
b/hudi-flink-datasource/hudi-flink/src/test/java/org/apache/hudi/table/catalog/TestHoodieHiveCatalog.java
@@ -43,7 +43,6 @@ import org.apache.hudi.utils.CatalogUtils;
 
 import com.fasterxml.jackson.databind.JsonNode;
 import com.fasterxml.jackson.databind.ObjectMapper;
-import org.apache.flink.calcite.shaded.com.google.common.collect.Lists;
 import org.apache.flink.table.api.DataTypes;
 import org.apache.flink.table.api.Schema;
 import org.apache.flink.table.catalog.AbstractCatalog;
@@ -71,6 +70,7 @@ import org.junit.jupiter.params.provider.ValueSource;
 
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
@@ -126,7 +126,7 @@ public class TestHoodieHiveCatalog extends 
BaseTestHoodieCatalog {
           .column("par2", DataTypes.STRING())
           .primaryKey("uuid")
           .build();
-  List<String> multiPartitions = Lists.newArrayList("par1", "par2");
+  List<String> multiPartitions = Arrays.asList("par1", "par2");
 
   private static HoodieHiveCatalog hoodieCatalog;
   private final ObjectPath tablePath = new ObjectPath(TEST_DEFAULT_DATABASE, 
"test");
@@ -276,7 +276,7 @@ public class TestHoodieHiveCatalog extends 
BaseTestHoodieCatalog {
     assertEquals(keyGeneratorClassName, 
NonpartitionedAvroKeyGenerator.class.getName());
 
     // validate the order of partition fields in the multi-partition table
-    List<String> multiPartitions =  Lists.newArrayList("par2", "par1");
+    List<String> multiPartitions = Arrays.asList("par2", "par1");
     ObjectPath multiPartitionsTablePath = new ObjectPath("default", "tb_mp_" + 
System.currentTimeMillis());
     CatalogTable multiPartitionsTable =
         CatalogUtils.createCatalogTable(singleKeyMultiPartitionTableSchema, 
multiPartitions, options, "multi-partition hudi table");
diff --git a/hudi-flink-datasource/hudi-flink1.18.x/pom.xml 
b/hudi-flink-datasource/hudi-flink1.18.x/pom.xml
index 99740ede3b62..ae8610a23913 100644
--- a/hudi-flink-datasource/hudi-flink1.18.x/pom.xml
+++ b/hudi-flink-datasource/hudi-flink1.18.x/pom.xml
@@ -127,12 +127,6 @@
             <version>${flink1.18.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-table-planner_2.12</artifactId>
-            <version>${flink1.18.version}</version>
-            <scope>provided</scope>
-        </dependency>
 
         <!-- Test dependencies -->
         <dependency>
diff --git a/hudi-flink-datasource/hudi-flink1.19.x/pom.xml 
b/hudi-flink-datasource/hudi-flink1.19.x/pom.xml
index 20b05d700ddb..bc70bb5e96b1 100644
--- a/hudi-flink-datasource/hudi-flink1.19.x/pom.xml
+++ b/hudi-flink-datasource/hudi-flink1.19.x/pom.xml
@@ -127,12 +127,6 @@
             <version>${flink1.19.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-table-planner_2.12</artifactId>
-            <version>${flink1.19.version}</version>
-            <scope>provided</scope>
-        </dependency>
 
         <!-- Test dependencies -->
         <dependency>
diff --git a/hudi-flink-datasource/hudi-flink1.20.x/pom.xml 
b/hudi-flink-datasource/hudi-flink1.20.x/pom.xml
index 6be0c90abf5d..266d472efade 100644
--- a/hudi-flink-datasource/hudi-flink1.20.x/pom.xml
+++ b/hudi-flink-datasource/hudi-flink1.20.x/pom.xml
@@ -127,12 +127,6 @@
             <version>${flink1.20.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-table-planner_2.12</artifactId>
-            <version>${flink1.20.version}</version>
-            <scope>provided</scope>
-        </dependency>
 
         <!-- Test dependencies -->
         <dependency>
diff --git a/hudi-flink-datasource/hudi-flink2.0.x/pom.xml 
b/hudi-flink-datasource/hudi-flink2.0.x/pom.xml
index c56d57618663..770b92a94bc1 100644
--- a/hudi-flink-datasource/hudi-flink2.0.x/pom.xml
+++ b/hudi-flink-datasource/hudi-flink2.0.x/pom.xml
@@ -127,12 +127,6 @@
             <version>${flink2.0.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-table-planner_2.12</artifactId>
-            <version>${flink2.0.version}</version>
-            <scope>provided</scope>
-        </dependency>
 
         <!-- Test dependencies -->
         <dependency>
diff --git a/hudi-flink-datasource/hudi-flink2.1.x/pom.xml 
b/hudi-flink-datasource/hudi-flink2.1.x/pom.xml
index 7304d4375bc2..5f3db5162bfd 100644
--- a/hudi-flink-datasource/hudi-flink2.1.x/pom.xml
+++ b/hudi-flink-datasource/hudi-flink2.1.x/pom.xml
@@ -121,12 +121,6 @@
             <version>${flink2.1.version}</version>
             <scope>provided</scope>
         </dependency>
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-table-planner_2.12</artifactId>
-            <version>${flink2.1.version}</version>
-            <scope>provided</scope>
-        </dependency>
 
         <!-- Test dependencies -->
         <dependency>
diff --git a/pom.xml b/pom.xml
index 292fd806d163..2523a3c7b5f6 100644
--- a/pom.xml
+++ b/pom.xml
@@ -2534,14 +2534,11 @@
                         <exclude>*:*_2.12</exclude>
                       </excludes>
                       <!--
-                        Flink decoupled from Scala starting 1.15; the `_2.12` 
suffix on
-                        flink-table-planner and its transitives (scala-xml, 
chill) is a
-                        legacy naming artifact, not a real Scala 2.12 binary 
dependency.
-                        Flink publishes no `_2.13` variant of 
flink-table-planner, so
-                        building with `-Dscala-2.13` together with any Flink 
profile
-                        (e.g. `-Dflink1.20`) would otherwise trip the blanket 
`*:*_2.12`
-                        ban above. Whitelist only the Flink-side transitives 
that
-                        actually surface here.
+                        Some Flink artifacts still carry a `_2.12` suffix even 
though
+                        newer Flink modules are decoupled from Scala. Building 
with
+                        `-Dscala-2.13` together with a Flink profile would 
otherwise
+                        trip the blanket `*:*_2.12` ban above. Whitelist only 
the
+                        Flink-side transitives that actually surface here.
                       -->
                       <includes>
                         <include>org.apache.flink:*_2.12</include>

Reply via email to