This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new ac8792761b [hive] Unescape partition path values when building Hive
split predicate (#9080)
ac8792761b is described below
commit ac8792761bbfca9f1f7359ac5a18ffb6bc1da1a4
Author: Eunbin Son <[email protected]>
AuthorDate: Thu Aug 6 23:12:25 2026 +0900
[hive] Unescape partition path values when building Hive split predicate
(#9080)
---
.../paimon/hive/utils/HiveSplitGenerator.java | 20 +++---
.../apache/paimon/hive/HiveSplitGeneratorTest.java | 74 ++++++++++++++++++++++
2 files changed, 83 insertions(+), 11 deletions(-)
diff --git
a/paimon-hive/paimon-hive-connector-common/src/main/java/org/apache/paimon/hive/utils/HiveSplitGenerator.java
b/paimon-hive/paimon-hive-connector-common/src/main/java/org/apache/paimon/hive/utils/HiveSplitGenerator.java
index 0117434296..a6a55f0aa7 100644
---
a/paimon-hive/paimon-hive-connector-common/src/main/java/org/apache/paimon/hive/utils/HiveSplitGenerator.java
+++
b/paimon-hive/paimon-hive-connector-common/src/main/java/org/apache/paimon/hive/utils/HiveSplitGenerator.java
@@ -18,6 +18,7 @@
package org.apache.paimon.hive.utils;
+import org.apache.paimon.fs.Path;
import org.apache.paimon.hive.HiveConnectorOptions;
import org.apache.paimon.hive.mapred.PaimonInputSplit;
import org.apache.paimon.io.DataFileMeta;
@@ -31,6 +32,7 @@ import org.apache.paimon.table.source.InnerTableScan;
import org.apache.paimon.tag.TagPreview;
import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.BinPacking;
+import org.apache.paimon.utils.PartitionPathUtils;
import org.apache.hadoop.hive.conf.HiveConf;
import org.apache.hadoop.mapred.InputSplit;
@@ -135,17 +137,13 @@ public class HiveSplitGenerator {
String defaultPartName) {
Set<String> partitionKeySet = new HashSet<>(partitionKeys);
LinkedHashMap<String, String> partition = new LinkedHashMap<>();
- for (String s : partitionDir.split("/")) {
- s = s.trim();
- if (s.isEmpty()) {
- continue;
- }
- String[] kv = s.split("=");
- if (kv.length != 2) {
- continue;
- }
- if (partitionKeySet.contains(kv[0])) {
- partition.put(kv[0], kv[1]);
+ // the directory names are escaped by
PartitionPathUtils#escapePathName when the partition
+ // is created, so they must be unescaped to get back the raw partition
values
+ LinkedHashMap<String, String> spec =
+ PartitionPathUtils.extractPartitionSpecFromPath(new
Path(partitionDir));
+ for (Map.Entry<String, String> entry : spec.entrySet()) {
+ if (partitionKeySet.contains(entry.getKey())) {
+ partition.put(entry.getKey(), entry.getValue());
}
}
if (partition.isEmpty() || partition.size() != partitionKeys.size()) {
diff --git
a/paimon-hive/paimon-hive-connector-common/src/test/java/org/apache/paimon/hive/HiveSplitGeneratorTest.java
b/paimon-hive/paimon-hive-connector-common/src/test/java/org/apache/paimon/hive/HiveSplitGeneratorTest.java
index 2c3fee3749..16aa824c1b 100644
---
a/paimon-hive/paimon-hive-connector-common/src/test/java/org/apache/paimon/hive/HiveSplitGeneratorTest.java
+++
b/paimon-hive/paimon-hive-connector-common/src/test/java/org/apache/paimon/hive/HiveSplitGeneratorTest.java
@@ -22,6 +22,8 @@ import org.apache.paimon.CoreOptions;
import org.apache.paimon.catalog.CatalogContext;
import org.apache.paimon.data.BinaryRow;
import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.data.GenericRow;
import org.apache.paimon.fs.FileIO;
import org.apache.paimon.fs.Path;
import org.apache.paimon.hive.utils.HiveSplitGenerator;
@@ -33,6 +35,9 @@ import org.apache.paimon.schema.TableSchema;
import org.apache.paimon.table.AppendOnlyFileStoreTable;
import org.apache.paimon.table.CatalogEnvironment;
import org.apache.paimon.table.FileStoreTable;
+import org.apache.paimon.table.sink.BatchTableCommit;
+import org.apache.paimon.table.sink.BatchTableWrite;
+import org.apache.paimon.table.sink.BatchWriteBuilder;
import org.apache.paimon.table.source.DataSplit;
import org.apache.paimon.types.DataField;
import org.apache.paimon.types.IntType;
@@ -40,7 +45,9 @@ import org.apache.paimon.types.VarCharType;
import org.apache.paimon.utils.TraceableFileIO;
import org.apache.hadoop.hive.conf.HiveConf;
+import org.apache.hadoop.mapred.InputSplit;
import org.apache.hadoop.mapred.JobConf;
+import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
@@ -72,6 +79,26 @@ public class HiveSplitGeneratorTest {
Collections.emptyMap(),
"");
+ private static final TableSchema PARTITIONED_TABLE_SCHEMA =
+ new TableSchema(
+ 0,
+ SCHEMA_FIELDS,
+ 2,
+ Collections.singletonList("pt"),
+ Collections.emptyList(),
+ Collections.emptyMap(),
+ "");
+
+ /** Partition value holding a character that {@code escapePathName}
rewrites. */
+ private static final String ESCAPED_PARTITION_VALUE = "2026-01-01 00:00";
+
+ /** Directory name Paimon writes for {@link #ESCAPED_PARTITION_VALUE}. */
+ private static final String ESCAPED_PARTITION_DIR = "pt=2026-01-01
00%3A00";
+
+ private static final String PLAIN_PARTITION_VALUE = "plain";
+
+ private static final String PLAIN_PARTITION_DIR = "pt=plain";
+
@TempDir java.nio.file.Path tempDir;
protected Path tablePath;
@@ -133,6 +160,53 @@ public class HiveSplitGeneratorTest {
assertThat(totalFiles).isEqualTo(10);
}
+ @Test
+ public void testGenerateSplitsForEscapedPartitionValue() throws Exception {
+ FileStoreTable table = createPartitionedTableWithTwoPartitions();
+
+ // Hive hands back the escaped directory registered in the metastore
+ assertThat(fileIO.exists(new Path(tablePath,
ESCAPED_PARTITION_DIR))).isTrue();
+
+ InputSplit[] splits = generateSplits(table, ESCAPED_PARTITION_DIR);
+ assertThat(splits).hasSize(1);
+ }
+
+ @Test
+ public void testGenerateSplitsForPlainPartitionValue() throws Exception {
+ FileStoreTable table = createPartitionedTableWithTwoPartitions();
+
+ assertThat(fileIO.exists(new Path(tablePath,
PLAIN_PARTITION_DIR))).isTrue();
+
+ InputSplit[] splits = generateSplits(table, PLAIN_PARTITION_DIR);
+ assertThat(splits).hasSize(1);
+ }
+
+ private InputSplit[] generateSplits(FileStoreTable table, String
partitionDir) {
+ JobConf jobConf = new JobConf();
+ jobConf.set(FileInputFormat.INPUT_DIR, tablePath + "/" + partitionDir);
+ return HiveSplitGenerator.generateSplits(table, jobConf, 1);
+ }
+
+ private FileStoreTable createPartitionedTableWithTwoPartitions() throws
Exception {
+ FileStoreTable table = createFileStoreTable(PARTITIONED_TABLE_SCHEMA);
+ BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
+ try (BatchTableWrite write = writeBuilder.newWrite();
+ BatchTableCommit commit = writeBuilder.newCommit()) {
+ write.write(
+ GenericRow.of(
+ 1,
+ BinaryString.fromString("a"),
+ BinaryString.fromString(ESCAPED_PARTITION_VALUE)));
+ write.write(
+ GenericRow.of(
+ 2,
+ BinaryString.fromString("b"),
+ BinaryString.fromString(PLAIN_PARTITION_VALUE)));
+ commit.commit(write.prepareCommit());
+ }
+ return table;
+ }
+
private FileStoreTable createFileStoreTable(TableSchema tableSchema)
throws Exception {
SchemaManager schemaManager = new SchemaManager(fileIO, tablePath);
schemaManager.commit(tableSchema);