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 f17346a6574c386bcf18e47165c70a3dac03ead9 Author: Danny Chan <[email protected]> AuthorDate: Thu Jul 2 16:24:25 2026 +0800 fix: remove the dependency to flink-table-planner (#19131) (cherry picked from commit 727ded72b2522efc3e5bc1dc369bf82814d14d75) --- hudi-examples/hudi-examples-flink/pom.xml | 13 +- hudi-flink-datasource/hudi-flink/pom.xml | 13 +- .../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 +- .../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 +- 21 files changed, 649 insertions(+), 100 deletions(-) diff --git a/hudi-examples/hudi-examples-flink/pom.xml b/hudi-examples/hudi-examples-flink/pom.xml index 75a21c17053a..da370499262c 100644 --- a/hudi-examples/hudi-examples-flink/pom.xml +++ b/hudi-examples/hudi-examples-flink/pom.xml @@ -194,12 +194,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> @@ -371,6 +365,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 1c8f323d0800..1923534f475b 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/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 caffecbb34ba..e1f03fdfb2a8 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.typeutils.RowDataSerializer; import org.apache.flink.table.runtime.generated.GeneratedNormalizedKeyComputer; import org.apache.flink.table.runtime.generated.GeneratedRecordComparator; @@ -127,9 +126,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 f67be66b3c3e..6e6695798825 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; @@ -322,8 +321,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( @@ -345,10 +345,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 036f724b00d5..8d39fae75092 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 @@ -79,7 +79,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; @@ -154,7 +153,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)) { @@ -184,7 +183,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); } } @@ -563,7 +562,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)) @@ -606,7 +605,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 107e8673c42b..c1b76db677ab 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/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 7f67a0281e5c..6f778b9bc297 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 4df05d62af7a..72adebafa4b1 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 5c29d042577e..983dde06bd55 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 fee65624c685..70838282e5ac 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 90c40c1f8959..322779f7c119 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 f6665d3c42b3..f8d07669b2a8 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 35ac14994c17..1d5d6e17ebfd 100644 --- a/pom.xml +++ b/pom.xml @@ -2514,14 +2514,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>
