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 c80f8a6857 [flink] Fix scan partitions with partial partition spec
(#9225)
c80f8a6857 is described below
commit c80f8a68577332186c8c5664a2d63bd358331326
Author: Arnav Balyan <[email protected]>
AuthorDate: Sat Aug 15 17:43:18 2026 +0530
[flink] Fix scan partitions with partial partition spec (#9225)
---
.../org/apache/paimon/flink/lookup/StaticPartitionLoader.java | 6 +++++-
.../java/org/apache/paimon/flink/BatchFileStoreITCase.java | 10 ++++++++++
2 files changed, 15 insertions(+), 1 deletion(-)
diff --git
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/StaticPartitionLoader.java
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/StaticPartitionLoader.java
index bdfecc39f9..1a7eb32abe 100644
---
a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/StaticPartitionLoader.java
+++
b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/lookup/StaticPartitionLoader.java
@@ -26,6 +26,7 @@ import org.apache.paimon.types.RowType;
import org.apache.paimon.utils.InternalRowPartitionComputer;
import java.util.ArrayList;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -48,9 +49,12 @@ public class StaticPartitionLoader extends PartitionLoader {
String defaultPartitionName =
Options.fromMap(table.options()).get(PARTITION_DEFAULT_NAME);
InternalRowSerializer serializer = new
InternalRowSerializer(partitionType);
for (Map<String, String> spec : scanPartitions) {
+ Map<String, String> completedSpec = new LinkedHashMap<>(spec);
+ table.partitionKeys()
+ .forEach(key -> completedSpec.putIfAbsent(key,
defaultPartitionName));
GenericRow row =
InternalRowPartitionComputer.convertSpecToInternalRow(
- spec, partitionType, defaultPartitionName);
+ completedSpec, partitionType,
defaultPartitionName);
partitions.add(serializer.toBinaryRow(row).copy());
}
}
diff --git
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
index 6c9570dddc..94f223f3de 100644
---
a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
+++
b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/BatchFileStoreITCase.java
@@ -1101,6 +1101,16 @@ public class BatchFileStoreITCase extends
CatalogITCaseBase {
assertThat(sql(query)).containsExactly(Row.of(1, 11), Row.of(1, 12),
Row.of(2, 22));
}
+ @Test
+ public void testScanWithPartialSpecifiedPartition() {
+ sql("CREATE TABLE P (dt STRING, hh INT, v INT) PARTITIONED BY (dt,
hh)");
+ sql(
+ "INSERT INTO P VALUES ('20260814', CAST(NULL AS INT), 1),
('20260814', 10, 2), ('20260815', CAST(NULL AS INT), 3)");
+
+ assertThat(sql("SELECT COUNT(*) FROM P /*+ OPTIONS('scan.partitions' =
'dt=20260814') */"))
+ .containsExactly(Row.of(1L));
+ }
+
@Test
public void testScanWithSpecifiedPartitionsWithFieldMapping() {
sql("CREATE TABLE P (id INT, v INT, pt STRING) PARTITIONED BY (pt)");