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 c55209ab1d [core] Support arbitrary time granularity for chain table
delta computation (#8185)
c55209ab1d is described below
commit c55209ab1d86a77989d4c492fb4296ee66720619
Author: Juntao Zhang <[email protected]>
AuthorDate: Sat Aug 15 14:37:27 2026 +0800
[core] Support arbitrary time granularity for chain table delta computation
(#8185)
---
docs/docs/primary-key-table/chain-table.mdx | 3 +
.../paimon/partition/PartitionTimeResolver.java | 443 +++++++++++++++++++++
.../org/apache/paimon/utils/ChainTableUtils.java | 135 ++-----
.../partition/PartitionTimeResolverTest.java | 374 +++++++++++++++++
.../apache/paimon/utils/ChainTableUtilsTest.java | 148 +++++--
.../apache/paimon/spark/SparkChainTableITCase.java | 57 +++
6 files changed, 1031 insertions(+), 129 deletions(-)
diff --git a/docs/docs/primary-key-table/chain-table.mdx
b/docs/docs/primary-key-table/chain-table.mdx
index 38e5be64de..1a3dfd1a63 100644
--- a/docs/docs/primary-key-table/chain-table.mdx
+++ b/docs/docs/primary-key-table/chain-table.mdx
@@ -136,6 +136,9 @@ Notice that:
to use streaming read or lookup join. Other merge engine types are not
supported
on the delta branch for these incremental read paths. Batch read is not
affected.
- Chain table requires `sequence.field`, so the `FIRST_ROW` and `AGGREGATE`
merge engines are not supported.
+- For chain table delta computation, `partition.timestamp-formatter` must
specify a complete date
+ (year-month-day or year-day-of-year), only use round-trippable date/time
fields
+ (`y/u`, `M/L`, `d`, `D`, `H`, `k`, `m`, `s`), and the minimum step must be
at least one second.
## Write Data
diff --git
a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolver.java
b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolver.java
new file mode 100644
index 0000000000..08cc4b4e97
--- /dev/null
+++
b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolver.java
@@ -0,0 +1,443 @@
+/*
+ * 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.paimon.partition;
+
+import java.text.ParsePosition;
+import java.time.DateTimeException;
+import java.time.Duration;
+import java.time.LocalDateTime;
+import java.time.Period;
+import java.time.format.DateTimeFormatter;
+import java.time.temporal.ChronoField;
+import java.time.temporal.ChronoUnit;
+import java.time.temporal.TemporalAccessor;
+import java.time.temporal.TemporalAmount;
+import java.time.temporal.TemporalField;
+import java.time.temporal.TemporalUnit;
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.Map;
+import java.util.Set;
+import java.util.stream.Collectors;
+
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+
+/**
+ * Resolves timestamp pattern and formatter to extract time step and compute
partition values for
+ * chain table partitions.
+ */
+public class PartitionTimeResolver {
+ private static final Map<Character, TemporalField> FIELD_MAP = new
HashMap<>();
+ private final PartitionTimeExtractor timeExtractor;
+ private final List<String> partitionKeys;
+ private final String pattern;
+ private final String formatter;
+ private Map<PatternToken, List<FormatToken>> patternFormatMappings;
+ private List<PatternToken> patternTokens;
+ private List<FormatToken> formatTokens;
+
+ public PartitionTimeResolver(List<String> partitionKeys, String pattern,
String formatter) {
+ checkArgument(pattern != null, "pattern cannot be null");
+ checkArgument(formatter != null, "formatter cannot be null");
+ checkArgument(partitionKeys != null, "partitionColumns cannot be
null");
+ this.timeExtractor = new PartitionTimeExtractor(pattern, formatter);
+ this.partitionKeys = partitionKeys;
+ this.pattern = pattern;
+ this.formatter = formatter;
+ init();
+ }
+
+ static {
+ // Only round-trippable date/time fields are supported for chain table
delta computation.
+ FIELD_MAP.put('y', ChronoField.YEAR);
+ FIELD_MAP.put('u', ChronoField.YEAR);
+ FIELD_MAP.put('M', ChronoField.MONTH_OF_YEAR);
+ FIELD_MAP.put('L', ChronoField.MONTH_OF_YEAR);
+ FIELD_MAP.put('d', ChronoField.DAY_OF_MONTH);
+ FIELD_MAP.put('D', ChronoField.DAY_OF_YEAR);
+ FIELD_MAP.put('H', ChronoField.HOUR_OF_DAY);
+ FIELD_MAP.put('k', ChronoField.CLOCK_HOUR_OF_DAY);
+ FIELD_MAP.put('m', ChronoField.MINUTE_OF_HOUR);
+ FIELD_MAP.put('s', ChronoField.SECOND_OF_MINUTE);
+ }
+
+ private void init() {
+ this.patternFormatMappings = new HashMap<>();
+ this.patternTokens = parsePattern();
+ this.formatTokens = parseFormatter();
+ boolean matched = matchRecursive(0, 0);
+ checkArgument(
+ matched, "Failed to match pattern '%s' to formatter '%s'",
pattern, formatter);
+ validateFormatterFields();
+ }
+
+ /**
+ * Validates that the formatter specifies a complete date. Unsupported
field letters have
+ * already been rejected by {@link #parseFormatter()}.
+ */
+ private void validateFormatterFields() {
+ Set<TemporalField> dateFields = new HashSet<>();
+ for (FormatToken token : formatTokens) {
+ if (token instanceof TemporalFieldToken) {
+ dateFields.add(((TemporalFieldToken) token).field);
+ }
+ }
+
+ boolean hasCompleteDate =
+ (dateFields.contains(ChronoField.YEAR)
+ &&
dateFields.contains(ChronoField.MONTH_OF_YEAR)
+ &&
dateFields.contains(ChronoField.DAY_OF_MONTH))
+ || (dateFields.contains(ChronoField.YEAR)
+ &&
dateFields.contains(ChronoField.DAY_OF_YEAR));
+ checkArgument(
+ hasCompleteDate,
+ "Formatter '%s' does not specify a complete date. "
+ + "Chain table delta computation requires
year-month-day or year-day-of-year.",
+ formatter);
+ }
+
+ /**
+ * Extracts the minimum time step from the given pattern and formatter.
+ *
+ * @return the smallest {@link Duration} or {@link Period} step among
variable-controlled time
+ * units
+ */
+ public TemporalAmount extractMinStep() {
+ TemporalAmount minStep = null;
+ Duration minDuration = null;
+ for (PatternToken patternToken : patternTokens) {
+ if (!patternToken.isVariable) {
+ continue;
+ }
+ List<FormatToken> tokens = patternFormatMappings.get(patternToken);
+ if (tokens == null || tokens.isEmpty()) {
+ continue;
+ }
+ for (FormatToken token : tokens) {
+ if (!(token instanceof TemporalFieldToken)) {
+ continue;
+ }
+ TemporalFieldToken fieldToken = (TemporalFieldToken) token;
+ Duration duration =
fieldToken.field.getBaseUnit().getDuration();
+ if (minDuration == null || duration.compareTo(minDuration) <
0) {
+ minDuration = duration;
+ minStep = stepOf(fieldToken);
+ }
+ }
+ }
+ checkArgument(minStep != null, "No time field found in pattern
variables");
+ return minStep;
+ }
+
+ /**
+ * Computes partition column values by formatting the given datetime and
extracting each
+ * variable's segment according to the pattern-to-format mapping.
+ */
+ public LinkedHashMap<String, String> resolvePartitionValues(LocalDateTime
dateTime) {
+ LinkedHashMap<String, String> result = new LinkedHashMap<>();
+ for (PatternToken patternToken : patternTokens) {
+ if (!patternToken.isVariable) {
+ continue;
+ }
+ String variableName = patternToken.token.substring(1);
+ List<FormatToken> tokens = patternFormatMappings.get(patternToken);
+ int start = tokens.get(0).start;
+ int end = tokens.get(tokens.size() - 1).end;
+ DateTimeFormatter fmt =
+ DateTimeFormatter.ofPattern(formatter.substring(start,
end), Locale.ROOT);
+ result.put(variableName, fmt.format(dateTime));
+ }
+ return result;
+ }
+
+ public LocalDateTime parsePartitionValues(List<?> partitionValues) {
+ return timeExtractor.extract(this.partitionKeys, partitionValues);
+ }
+
+ /** Parses formatter into format tokens (time fields and literals). */
+ private List<FormatToken> parseFormatter() {
+ List<FormatToken> tokens = new ArrayList<>();
+ for (int pos = 0; pos < formatter.length(); pos++) {
+ char c = formatter.charAt(pos);
+ if (isTimeChar(c)) {
+ int start = pos;
+ while (pos < formatter.length() && formatter.charAt(pos) == c)
{
+ pos++;
+ }
+ TemporalField field = FIELD_MAP.get(c);
+ tokens.add(new TemporalFieldToken(c, field, start, pos));
+ pos--;
+ } else if (c == '\'') {
+ // parse literals
+ int start = pos++;
+ for (; pos < formatter.length(); pos++) {
+ if (formatter.charAt(pos) == '\'') {
+ if (pos + 1 < formatter.length() &&
formatter.charAt(pos + 1) == '\'') {
+ pos++;
+ } else {
+ break; // end of literal
+ }
+ }
+ }
+ checkArgument(
+ pos < formatter.length(),
+ "Pattern ends with an incomplete string literal: " +
formatter);
+ String str = formatter.substring(start + 1, pos);
+ if (str.isEmpty()) {
+ tokens.add(new LiteralToken("'", start, pos + 1));
+ } else {
+ tokens.add(new LiteralToken(str.replace("''", "'"), start,
pos + 1));
+ }
+ } else if (Character.isLetter(c)) {
+ throw new IllegalArgumentException(
+ String.format(
+ "Unsupported formatter pattern letter '%s' in
formatter: %s.",
+ c, formatter));
+ } else {
+ tokens.add(new LiteralToken(String.valueOf(c), pos, pos + 1));
+ }
+ }
+ checkArgument(!tokens.isEmpty(), "No time unit found in formatter:
%s", formatter);
+ return tokens;
+ }
+
+ private static boolean isTimeChar(char c) {
+ return FIELD_MAP.containsKey(c);
+ }
+
+ /** Parses pattern string into pattern tokens (variables and literals). */
+ private List<PatternToken> parsePattern() {
+ List<String> sortedPartCols =
+ partitionKeys.stream()
+ .sorted(Comparator.reverseOrder())
+ .collect(Collectors.toList());
+
+ List<PatternToken> tokens = new ArrayList<>();
+ StringBuilder literalBuf = new StringBuilder();
+ for (int cursor = 0, len = pattern.length(); cursor < len; ) {
+ char curr = pattern.charAt(cursor);
+ if (curr == '$') {
+ if (literalBuf.length() > 0) {
+ tokens.add(new PatternToken(literalBuf.toString(), false));
+ literalBuf.setLength(0);
+ }
+ boolean matched = false;
+ // Match the longest column name first to resolve ambiguity
when one column name
+ // is a prefix of another (e.g., "dt" vs "dt1").
+ for (String part : sortedPartCols) {
+ String varToken = curr + part;
+ if (pattern.startsWith(varToken, cursor)) {
+ tokens.add(new PatternToken(varToken, true));
+ cursor += varToken.length();
+ matched = true;
+ break;
+ }
+ }
+ checkArgument(
+ matched,
+ "Unknown variable in pattern '%s' at position %s",
+ pattern,
+ cursor);
+ } else {
+ literalBuf.append(curr);
+ cursor++;
+ }
+ }
+ if (literalBuf.length() > 0) {
+ tokens.add(new PatternToken(literalBuf.toString(), false));
+ }
+ return tokens;
+ }
+
+ /**
+ * Recursively matches pattern tokens to format tokens. For variable
tokens, greedily consumes
+ * consecutive format tokens. For literal tokens, verifies length and
content match.
+ */
+ private boolean matchRecursive(int patternIdx, int formatIdx) {
+ if (patternIdx == patternTokens.size()) {
+ return formatIdx == formatTokens.size();
+ }
+
+ // Remaining format tokens must be at least as many as remaining
pattern tokens
+ if (formatTokens.size() - formatIdx < patternTokens.size() -
patternIdx) {
+ return false;
+ }
+
+ PatternToken patternToken = patternTokens.get(patternIdx);
+ // Max format tokens this pattern token can consume, leaving at least
1 token per remaining
+ // pattern token
+ int maxLen = formatTokens.size() - formatIdx - (patternTokens.size() -
patternIdx - 1);
+
+ int matchedEndIdx = -1;
+ for (int len = 1; len <= maxLen; len++) {
+ int formatEndIdx = formatIdx + len;
+ if (patternToken.isVariable) {
+ if (matchRecursive(patternIdx + 1, formatEndIdx)) {
+ checkArgument(
+ matchedEndIdx == -1,
+ "Ambiguous mapping for pattern variable '%s' in
pattern '%s' with formatter '%s'. "
+ + "Please separate adjacent variables with
literals.",
+ patternToken.token,
+ pattern,
+ formatter);
+ matchedEndIdx = formatEndIdx;
+ }
+ } else {
+ // Literal pattern tokens match 1...len consecutive format
tokens, split by token
+ // length
+ if (matchLiteral(patternToken.token, formatIdx, formatEndIdx))
{
+ if (matchRecursive(patternIdx + 1, formatEndIdx)) {
+ return true;
+ }
+ }
+ }
+ }
+ if (matchedEndIdx != -1) {
+ patternFormatMappings.put(patternToken,
formatTokens.subList(formatIdx, matchedEndIdx));
+ return true;
+ }
+ return false;
+ }
+
+ /** Checks if a literal pattern token matches a sequence of format tokens.
*/
+ private boolean matchLiteral(String literalToken, int startIdx, int
endIdx) {
+ StringBuilder subFormatter = new StringBuilder();
+ StringBuilder literalValue = new StringBuilder();
+ boolean pureLiteral = true;
+ List<TemporalField> fields = new ArrayList<>();
+ for (int i = startIdx; i < endIdx; i++) {
+ FormatToken token = formatTokens.get(i);
+ subFormatter.append(formatter, token.start, token.end);
+ pureLiteral = pureLiteral && token instanceof LiteralToken;
+ if (token instanceof TemporalFieldToken) {
+ fields.add(((TemporalFieldToken) (token)).field);
+ }
+ if (pureLiteral) {
+ literalValue.append(((LiteralToken) token).token);
+ }
+ }
+
+ if (pureLiteral) {
+ return literalToken.contentEquals(literalValue);
+ }
+
+ DateTimeFormatter fmt =
DateTimeFormatter.ofPattern(subFormatter.toString(), Locale.ROOT);
+ ParsePosition pp = new ParsePosition(0);
+ try {
+ TemporalAccessor ta = fmt.parse(literalToken, pp);
+ if (pp.getErrorIndex() >= 0 || pp.getIndex() !=
literalToken.length()) {
+ return false;
+ }
+ for (TemporalField field : fields) {
+ if (ta.isSupported(field)) {
+ try {
+ ta.get(field);
+ } catch (DateTimeException ignored) {
+ return false;
+ }
+ }
+ }
+ } catch (Exception ignored) {
+ return false;
+ }
+ return true;
+ }
+
+ private static TemporalAmount stepOf(TemporalFieldToken fieldToken) {
+ TemporalUnit unit = fieldToken.field.getBaseUnit();
+ if (unit == ChronoUnit.YEARS) {
+ return Period.ofYears(1);
+ }
+ if (unit == ChronoUnit.MONTHS) {
+ return Period.ofMonths(1);
+ }
+ return unit.getDuration();
+ }
+
+ private static class FormatToken {
+ final int start;
+ final int end;
+
+ private FormatToken(int start, int end) {
+ this.start = start;
+ this.end = end;
+ }
+
+ public int getLength() {
+ return end - start;
+ }
+ }
+
+ private static class LiteralToken extends FormatToken {
+ final String token;
+
+ LiteralToken(String token, int start, int end) {
+ super(start, end);
+ this.token = token;
+ }
+
+ @Override
+ public int getLength() {
+ return token.length();
+ }
+
+ @Override
+ public String toString() {
+ return String.format("LiteralToken{token=%s, start=%d, end=%d}",
token, start, end);
+ }
+ }
+
+ private static class TemporalFieldToken extends FormatToken {
+ final char letter;
+ final TemporalField field;
+
+ TemporalFieldToken(char letter, TemporalField field, int start, int
end) {
+ super(start, end);
+ this.letter = letter;
+ this.field = field;
+ }
+
+ @Override
+ public String toString() {
+ return String.format(
+ "TimeFieldToken{letter=%s, field=%s, start=%d, end=%d}",
+ letter, field, start, end);
+ }
+ }
+
+ private static class PatternToken {
+ final String token;
+ final boolean isVariable;
+
+ PatternToken(String token, boolean isVariable) {
+ this.token = token;
+ this.isVariable = isVariable;
+ }
+
+ @Override
+ public String toString() {
+ return String.format("PatternToken{token='%s', isVariable=%s}",
token, isVariable);
+ }
+ }
+}
diff --git
a/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
b/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
index 4a7bb1b073..56d2850cea 100644
--- a/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
+++ b/paimon-core/src/main/java/org/apache/paimon/utils/ChainTableUtils.java
@@ -26,7 +26,7 @@ import
org.apache.paimon.data.serializer.InternalRowSerializer;
import org.apache.paimon.io.DataFileMeta;
import org.apache.paimon.mergetree.SortedRun;
import org.apache.paimon.mergetree.compact.IntervalPartition;
-import org.apache.paimon.partition.PartitionTimeExtractor;
+import org.apache.paimon.partition.PartitionTimeResolver;
import org.apache.paimon.predicate.Predicate;
import org.apache.paimon.predicate.PredicateBuilder;
import org.apache.paimon.table.ChainGroupReadTable;
@@ -39,8 +39,10 @@ import org.apache.paimon.types.RowType;
import javax.annotation.Nullable;
+import java.time.Duration;
import java.time.LocalDateTime;
-import java.time.format.DateTimeFormatter;
+import java.time.temporal.ChronoUnit;
+import java.time.temporal.TemporalAmount;
import java.util.ArrayList;
import java.util.Collections;
import java.util.Comparator;
@@ -52,8 +54,6 @@ import java.util.Optional;
import java.util.Set;
import java.util.function.BiFunction;
import java.util.function.Function;
-import java.util.regex.Matcher;
-import java.util.regex.Pattern;
import java.util.stream.Collectors;
import static org.apache.paimon.utils.Preconditions.checkArgument;
@@ -61,6 +61,7 @@ import static
org.apache.paimon.utils.Preconditions.checkNotNull;
/** Utils for chain table. */
public class ChainTableUtils {
+ public static final int MAX_DELTA_PARTITIONS = 10_000_000;
public static boolean isChainTable(Map<String, String> tblOptions) {
return CoreOptions.fromMap(tblOptions).isChainTable();
@@ -92,61 +93,50 @@ public class ChainTableUtils {
List<String> partitionColumns,
RowType partType,
CoreOptions options,
- RecordComparator partitionComparator,
InternalRowPartitionComputer partitionComputer) {
InternalRowSerializer serializer = new InternalRowSerializer(partType);
List<BinaryRow> deltaPartitions = new ArrayList<>();
- boolean isDailyPartition = partitionColumns.size() == 1;
List<String> startPartitionValues =
new
ArrayList<>(partitionComputer.generatePartValues(beginPartition).values());
List<String> endPartitionValues =
new
ArrayList<>(partitionComputer.generatePartValues(endPartition).values());
- PartitionTimeExtractor timeExtractor =
- new PartitionTimeExtractor(
- options.partitionTimestampPattern(),
options.partitionTimestampFormatter());
- LocalDateTime stratPartitionTime =
- timeExtractor.extract(partitionColumns, startPartitionValues);
- LocalDateTime candidateTime = stratPartitionTime;
- LocalDateTime endPartitionTime =
- timeExtractor.extract(partitionColumns, endPartitionValues);
+ PartitionTimeResolver timeResolver =
+ new PartitionTimeResolver(
+ partitionColumns,
+ options.partitionTimestampPattern(),
+ options.partitionTimestampFormatter());
+ LocalDateTime startPartitionTime =
timeResolver.parsePartitionValues(startPartitionValues);
+ LocalDateTime endPartitionTime =
timeResolver.parsePartitionValues(endPartitionValues);
+ TemporalAmount step = timeResolver.extractMinStep();
+ // Period-based steps (year/month) are coarse-grained and cannot
explode, so only
+ // fine-grained Duration steps need the partition-count guard.
+ if (step instanceof Duration) {
+ long totalSeconds = ChronoUnit.SECONDS.between(startPartitionTime,
endPartitionTime);
+ long stepSeconds = ((Duration) step).getSeconds();
+ long estimatedCount = stepSeconds == 0 ? 0 : totalSeconds /
stepSeconds;
+ checkArgument(
+ estimatedCount < MAX_DELTA_PARTITIONS,
+ "Too many delta partitions generated between '%s' and '%s'
"
+ + "(exceeds %s). Please widen the partition
granularity or reduce "
+ + "the query range. Pattern: '%s', formatter:
'%s'.",
+ startPartitionValues,
+ endPartitionValues,
+ MAX_DELTA_PARTITIONS,
+ options.partitionTimestampPattern(),
+ options.partitionTimestampFormatter());
+ }
+ LocalDateTime candidateTime = startPartitionTime.plus(step);
while (!candidateTime.isAfter(endPartitionTime)) {
- if (isDailyPartition) {
- if (candidateTime.isAfter(stratPartitionTime)) {
- deltaPartitions.add(
- serializer
- .toBinaryRow(
-
InternalRowPartitionComputer.convertSpecToInternalRow(
- calPartValues(
- candidateTime,
- partitionColumns,
-
options.partitionTimestampPattern(),
-
options.partitionTimestampFormatter()),
- partType,
-
options.partitionDefaultName()))
- .copy());
- }
- } else {
- for (int hour = 0; hour <= 23; hour++) {
- candidateTime =
candidateTime.toLocalDate().atStartOfDay().plusHours(hour);
- BinaryRow candidatePartition =
- serializer
- .toBinaryRow(
-
InternalRowPartitionComputer.convertSpecToInternalRow(
- calPartValues(
- candidateTime,
- partitionColumns,
-
options.partitionTimestampPattern(),
-
options.partitionTimestampFormatter()),
- partType,
-
options.partitionDefaultName()))
- .copy();
- if (partitionComparator.compare(candidatePartition,
beginPartition) > 0
- && partitionComparator.compare(candidatePartition,
endPartition) <= 0) {
- deltaPartitions.add(candidatePartition);
- }
- }
- }
- candidateTime =
candidateTime.toLocalDate().plusDays(1).atStartOfDay();
+ BinaryRow candidatePartition =
+ serializer
+ .toBinaryRow(
+
InternalRowPartitionComputer.convertSpecToInternalRow(
+
timeResolver.resolvePartitionValues(candidateTime),
+ partType,
+ options.partitionDefaultName()))
+ .copy();
+ deltaPartitions.add(candidatePartition);
+ candidateTime = candidateTime.plus(step);
}
return deltaPartitions;
}
@@ -183,48 +173,6 @@ public class ChainTableUtils {
return PredicateBuilder.and(fieldPredicates);
}
- public static LinkedHashMap<String, String> calPartValues(
- LocalDateTime dateTime,
- List<String> partitionKeys,
- String timestampPattern,
- String timestampFormatter) {
- DateTimeFormatter formatter =
DateTimeFormatter.ofPattern(timestampFormatter);
- String formattedDateTime = dateTime.format(formatter);
- Pattern keyPattern = Pattern.compile("\\$(\\w+)");
- Matcher keyMatcher = keyPattern.matcher(timestampPattern);
- List<String> keyOrder = new ArrayList<>();
- StringBuilder regexBuilder = new StringBuilder();
- int lastPosition = 0;
- while (keyMatcher.find()) {
- regexBuilder.append(
- Pattern.quote(timestampPattern.substring(lastPosition,
keyMatcher.start())));
- regexBuilder.append("(.+)");
- keyOrder.add(keyMatcher.group(1));
- lastPosition = keyMatcher.end();
- }
-
regexBuilder.append(Pattern.quote(timestampPattern.substring(lastPosition)));
-
- Matcher valueMatcher =
Pattern.compile(regexBuilder.toString()).matcher(formattedDateTime);
- if (!valueMatcher.matches() || valueMatcher.groupCount() !=
keyOrder.size()) {
- throw new IllegalArgumentException(
- "Formatted datetime does not match timestamp pattern");
- }
-
- Map<String, String> keyValues = new HashMap<>();
- for (int i = 0; i < keyOrder.size(); i++) {
- keyValues.put(keyOrder.get(i), valueMatcher.group(i + 1));
- }
- List<String> values =
- partitionKeys.stream()
- .map(key -> keyValues.getOrDefault(key, ""))
- .collect(Collectors.toList());
- LinkedHashMap<String, String> res = new LinkedHashMap<>();
- for (int i = 0; i < partitionKeys.size(); i++) {
- res.put(partitionKeys.get(i), values.get(i));
- }
- return res;
- }
-
public static boolean isScanFallbackDeltaBranch(CoreOptions options) {
return options.isChainTable()
&&
options.scanFallbackDeltaBranch().equalsIgnoreCase(options.branch());
@@ -329,7 +277,6 @@ public class ChainTableUtils {
chainPartitionColumns,
chainPartType,
options,
- chainPartitionComparator,
chainPartitionComputer);
// Combine each chain-only BinaryRow with the group part into a full
partition
diff --git
a/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeResolverTest.java
b/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeResolverTest.java
new file mode 100644
index 0000000000..4a19613948
--- /dev/null
+++
b/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeResolverTest.java
@@ -0,0 +1,374 @@
+/*
+ * 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.paimon.partition;
+
+import org.apache.paimon.CoreOptions;
+import org.apache.paimon.codegen.RecordComparator;
+import org.apache.paimon.data.BinaryRow;
+import org.apache.paimon.data.BinaryRowWriter;
+import org.apache.paimon.data.BinaryString;
+import org.apache.paimon.options.Options;
+import org.apache.paimon.types.DataTypes;
+import org.apache.paimon.types.RowType;
+import org.apache.paimon.utils.ChainPartitionProjector;
+import org.apache.paimon.utils.ChainTableUtils;
+
+import org.apache.paimon.shade.guava30.com.google.common.collect.ImmutableMap;
+
+import org.assertj.core.util.Lists;
+import org.junit.jupiter.api.Test;
+
+import java.time.Duration;
+import java.time.LocalDateTime;
+import java.time.Period;
+import java.time.temporal.TemporalAmount;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/** Tests for {@link PartitionTimeResolver}. */
+public class PartitionTimeResolverTest {
+
+ private TemporalAmount extractMinStep(
+ String pattern, String formatter, String... partitionKeys) {
+ return new PartitionTimeResolver(Arrays.asList(partitionKeys),
pattern, formatter)
+ .extractMinStep();
+ }
+
+ /** Extract a string value from a BinaryRow at the given position. */
+ private static String getString(BinaryRow row, int pos) {
+ return row.getString(pos).toString();
+ }
+
+ private static BinaryRow row(List<String> values) {
+ BinaryRow row = new BinaryRow(values.size());
+ BinaryRowWriter writer = new BinaryRowWriter(row);
+ for (int i = 0; i < values.size(); i++) {
+ writer.writeString(i, BinaryString.fromString(values.get(i)));
+ }
+ writer.complete();
+ return row;
+ }
+
+ @Test
+ public void testExtractMinStep() {
+ assertThat(extractMinStep("$y$M$d$H$m$s", "yyyyMMddHHmmss", "y", "M",
"d", "H", "m", "s"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("$y$M$d $H$m$s", "yyyyMMdd HHmmss", "y",
"M", "d", "H", "m", "s"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(
+ extractMinStep(
+ "$y-$M-$d $H:$m:$s",
+ "yyyy-MM-dd HH:mm:ss",
+ "y",
+ "M",
+ "d",
+ "H",
+ "m",
+ "s"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(
+ extractMinStep(
+ "$y-$M-$d T $H:$m:$s",
+ "yyyy-MM-dd 'T' HH:mm:ss",
+ "y",
+ "M",
+ "d",
+ "H",
+ "m",
+ "s"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("$a", "yyyyMMddHHmmss",
"a")).isEqualTo(Duration.ofSeconds(1));
+
+ assertThat(
+ extractMinStep(
+ "$a $aaT$aaa $a4Z", "yyMM dd'T'HHmm ss'Z'",
"a4", "aa", "a", "aaa"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("$a12$aaT$aaa00Z", "yyyyMMdd'T'HHmmss'Z'",
"aa", "aaa", "a"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("$aT$a1$a200", "yyyyMMdd'T'HHmmss", "a",
"a1", "a2"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("$aT$aa", "yyyyMMdd'T'HHmm", "a", "aa"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("$a", "yyyyMMdd'T'HHmmss",
"a")).isEqualTo(Duration.ofSeconds(1));
+
+ assertThat(
+ extractMinStep(
+ "$ab $c $d:$e:$f", "yyyyMM dd HH:mm:ss", "ab",
"c", "d", "e", "f"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("$day $a:$b", "yyyyMMdd HH:mm", "day", "a",
"b"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("$aa $a", "yyyy/MM/dd HH", "aa", "a"))
+ .isEqualTo(Duration.ofHours(1));
+
+ assertThat(extractMinStep("$a $b", "HH:mm:ss yyyyMMdd", "a", "b"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("12:$a $b", "HH:mm:ss yyyyMMdd", "b", "a"))
+ .isEqualTo(Duration.ofSeconds(1));
+ assertThat(extractMinStep("12:$a:01 $b", "HH:mm:ss yyyyMMdd", "a",
"b"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("12:02:01 $b", "HH:mm:ss yyyyMMdd", "b"))
+ .isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("$hour:00:00 $date", "HH:mm:ss yyyyMMdd",
"date", "hour"))
+ .isEqualTo(Duration.ofHours(1));
+ assertThat(extractMinStep("00:00:00 $b", "HH:mm:ss yyyyMMdd", "b"))
+ .isEqualTo(Duration.ofDays(1));
+ assertThat(
+ extractMinStep(
+ "$hour_minute:01 $date",
+ "HH:mm:ss yyyyMMdd",
+ "hour_minute",
+ "date"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(extractMinStep("12:$a $b", "HH:mm:ss yyMMdd", "b", "a"))
+ .isEqualTo(Duration.ofSeconds(1));
+
+ // Unused partition columns should not affect the extracted minimum
step.
+ assertThat(extractMinStep("$dt", "yyyy-MM-dd", "other", "dt"))
+ .isEqualTo(Duration.ofDays(1));
+ assertThat(
+ extractMinStep(
+ "$dt $hour:$minute:00",
+ "yyyy-MM-dd HH:mm:ss",
+ "region",
+ "dt",
+ "hour",
+ "minute"))
+ .isEqualTo(Duration.ofMinutes(1));
+ assertThat(
+ extractMinStep(
+ "$hour:00:00 $date", "HH:mm:ss yyyyMMdd",
"date", "extra", "hour"))
+ .isEqualTo(Duration.ofHours(1));
+
+ assertThat(extractMinStep("$a-01-$b", "yyyy-MM-dd", "a", "b"))
+ .isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("$a-01", "yyyy-MM-dd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$y/$m/$d", "yyyy/MM/dd", "d", "y", "m"))
+ .isEqualTo(Duration.ofDays(1));
+
+ assertThat(extractMinStep("$a", "yyyyMMdd",
"a")).isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("$a01", "yyyyMMdd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$a $aa", "yyyyMM dd", "a",
"aa")).isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("202601$a", "yyyyMMdd",
"a")).isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("2026$a01", "yyyyMMdd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$a1201", "yyyyMMdd",
"a")).isEqualTo(Period.ofYears(1));
+ assertThat(extractMinStep("$a01", "yyyyMMdd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$a1201", "yyyyMMdd",
"a")).isEqualTo(Period.ofYears(1));
+
+ assertThat(extractMinStep("$a01", "yyMMdd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$a1", "yyMd",
"a")).isEqualTo(Period.ofMonths(1));
+ assertThat(extractMinStep("$a1201", "yyMMdd",
"a")).isEqualTo(Period.ofYears(1));
+ assertThat(extractMinStep("$a-12-1", "yy-M-d",
"a")).isEqualTo(Period.ofYears(1));
+ assertThat(extractMinStep("$a $aa", "yyMM dd", "a",
"aa")).isEqualTo(Duration.ofDays(1));
+ assertThat(extractMinStep("$a'", "yyMMdd''",
"a")).isEqualTo(Duration.ofDays(1));
+
+ assertThat(extractMinStep("$dt", "yyyy-DDD",
"dt")).isEqualTo(Duration.ofDays(1));
+
+ // 'k' = clock-hour-of-day
+ assertThat(extractMinStep("$dt", "yyyyMMddkk",
"dt")).isEqualTo(Duration.ofHours(1));
+ }
+
+ @Test
+ public void testResolvePartitionValues() {
+ Map<String, String> partitionValues =
+ new PartitionTimeResolver(
+ Arrays.asList("dt", "hour"), "$dt
$hour:00:00", "yyyyMMdd HH:mm:ss")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 1,
12, 0, 0));
+ assertEquals(ImmutableMap.of("dt", "20230101", "hour", "12"),
partitionValues);
+
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt", "hr"), "$dt
$hr", "yyyyMMdd HH")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 2,
3, 0, 0));
+ assertEquals(ImmutableMap.of("dt", "20230102", "hr", "03"),
partitionValues);
+
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt"), "$dt",
"yyyyMMdd")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 1,
0, 0, 0));
+ assertEquals(ImmutableMap.of("dt", "20230101"), partitionValues);
+
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt", "t"), "$dtT$t",
"yy-M-d'T'H:m:ss")
+ .resolvePartitionValues(LocalDateTime.of(2023, 12, 1,
11, 2, 3));
+ assertEquals(ImmutableMap.of("dt", "23-12-1", "t", "11:2:03"),
partitionValues);
+
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt"), "$dt",
"yy-MMM-d")
+ .resolvePartitionValues(LocalDateTime.of(2023, 12, 1,
11, 2, 3));
+ assertEquals(ImmutableMap.of("dt", "23-Dec-1"), partitionValues);
+
+ // Partition columns that are not referenced by the pattern should not
appear in the result.
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("other", "dt"), "$dt",
"yyyy-MM-dd")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 1,
0, 0, 0));
+ assertEquals(ImmutableMap.of("dt", "2023-01-01"), partitionValues);
+
+ partitionValues =
+ new PartitionTimeResolver(
+ Arrays.asList("region", "dt", "hour"),
+ "$dt $hour:00:00",
+ "yyyy-MM-dd HH:mm:ss")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 1,
10, 0, 0));
+ assertEquals(ImmutableMap.of("dt", "2023-01-01", "hour", "10"),
partitionValues);
+
+ // Day-of-year is also a valid complete date.
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt"), "$dt",
"yyyy-DDD")
+ .resolvePartitionValues(LocalDateTime.of(2026, 8, 10,
15, 30, 0));
+ assertEquals(ImmutableMap.of("dt", "2026-222"), partitionValues);
+
+ // 'u' = year, same as 'y'
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt"), "$dt",
"uuuuMMdd")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 20,
0, 0));
+ assertEquals(ImmutableMap.of("dt", "20230120"), partitionValues);
+
+ // 'L' = month-of-year, same as 'M'
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt"), "$dt",
"yyyy-L-dd")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 20,
0, 0));
+ assertEquals(ImmutableMap.of("dt", "2023-1-20"), partitionValues);
+
+ // 'k' = clock-hour-of-day
+ partitionValues =
+ new PartitionTimeResolver(Arrays.asList("dt", "hour"), "$dt
$hour", "yyyyMMdd kk")
+ .resolvePartitionValues(LocalDateTime.of(2023, 1, 1,
12, 0));
+ assertEquals(ImmutableMap.of("dt", "20230101", "hour", "12"),
partitionValues);
+ }
+
+ @Test
+ public void testParsePartitionValuesWithHourMinuteGranularity() {
+ // partition keys: (region, dt, hour_minute), chain keys: (dt,
hour_minute)
+ RowType fullType =
+ RowType.builder()
+ .field("region", DataTypes.STRING().notNull())
+ .field("dt", DataTypes.STRING().notNull())
+ .field("hour_minute", DataTypes.STRING().notNull())
+ .build();
+
+ ChainPartitionProjector projector = new
ChainPartitionProjector(fullType, 2);
+
+ // Compare chain partition (dt, hour_minute) lexicographically
+ RecordComparator chainComparator = (a, b) ->
a.getString(1).compareTo(b.getString(1));
+
+ Options opts = new Options();
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_PATTERN, "$dtT$hour_minute");
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_FORMATTER, "yyyyMMdd'T'HHmm");
+ CoreOptions options = new CoreOptions(opts);
+
+ BinaryRow begin = row(Lists.newArrayList("CN", "20260609", "1010"));
+ BinaryRow end = row(Lists.newArrayList("CN", "20260609", "1015"));
+
+ List<BinaryRow> deltas =
+ ChainTableUtils.getDeltaPartitionsWithProjector(
+ begin, end, options, chainComparator, projector);
+
+ assertThat(deltas).hasSize(5);
+ for (BinaryRow delta : deltas) {
+ assertThat(getString(delta, 0)).isEqualTo("CN");
+ assertThat(getString(delta, 1)).isEqualTo("20260609");
+ }
+ assertThat(getString(deltas.get(0), 2)).isEqualTo("1011");
+ assertThat(getString(deltas.get(1), 2)).isEqualTo("1012");
+ assertThat(getString(deltas.get(2), 2)).isEqualTo("1013");
+ assertThat(getString(deltas.get(3), 2)).isEqualTo("1014");
+ assertThat(getString(deltas.get(4), 2)).isEqualTo("1015");
+ }
+
+ @Test
+ public void testParsePartitionValuesWithSeparateHourAndMinute() {
+ // partition keys: (region, dt, hour, minute), chain keys: (dt, hour,
minute)
+ RowType fullType =
+ RowType.builder()
+ .field("region", DataTypes.STRING().notNull())
+ .field("dt", DataTypes.STRING().notNull())
+ .field("hour", DataTypes.STRING().notNull())
+ .field("minute", DataTypes.STRING().notNull())
+ .build();
+
+ ChainPartitionProjector projector = new
ChainPartitionProjector(fullType, 3);
+
+ // Compare chain partition (dt, hour, minute) lexicographically
+ RecordComparator chainComparator = (a, b) ->
a.getString(2).compareTo(b.getString(2));
+
+ Options opts = new Options();
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_PATTERN,
"$dtT$hour$minute00");
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_FORMATTER,
"yyyyMMdd'T'HHmmss");
+ CoreOptions options = new CoreOptions(opts);
+
+ BinaryRow begin = row(Lists.newArrayList("CN", "20260609", "10",
"10"));
+ BinaryRow end = row(Lists.newArrayList("CN", "20260609", "10", "15"));
+
+ List<BinaryRow> deltas =
+ ChainTableUtils.getDeltaPartitionsWithProjector(
+ begin, end, options, chainComparator, projector);
+
+ assertThat(deltas).hasSize(5);
+ for (BinaryRow delta : deltas) {
+ assertThat(getString(delta, 0)).isEqualTo("CN");
+ assertThat(getString(delta, 1)).isEqualTo("20260609");
+ assertThat(getString(delta, 2)).isEqualTo("10");
+ }
+ assertThat(getString(deltas.get(0), 3)).isEqualTo("11");
+ assertThat(getString(deltas.get(1), 3)).isEqualTo("12");
+ assertThat(getString(deltas.get(2), 3)).isEqualTo("13");
+ assertThat(getString(deltas.get(3), 3)).isEqualTo("14");
+ assertThat(getString(deltas.get(4), 3)).isEqualTo("15");
+ }
+
+ @Test
+ public void testUnsupportedOrIncompleteGranularity() {
+ // Incomplete date: year only / year-month / quarter.
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyy", "dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("does not specify a complete date");
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyyMM", "dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("does not specify a complete date");
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyy-'Q'q", "dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+
+ // Non-unique fields.
+ assertThatThrownBy(() -> extractMinStep("$dt $ampm", "yyyyMMdd a",
"dt", "ampm"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyy-MM-W", "dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+
+ // Week-based year.
+ assertThatThrownBy(() -> extractMinStep("$dt", "YYYY-'W'ww-e", "dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+
+ // Sub-second step.
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyyMMddHHmmss.SSS",
"dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+
+ // Zone offset currently is not support
+ assertThatThrownBy(() -> extractMinStep("$dt", "yyyyMMddHHmmssZ",
"dt"))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Unsupported formatter pattern letter");
+ }
+}
diff --git
a/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
index 50e74d4d4e..9ddcf7dfda 100644
--- a/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
+++ b/paimon-core/src/test/java/org/apache/paimon/utils/ChainTableUtilsTest.java
@@ -40,12 +40,10 @@ import org.assertj.core.util.Lists;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashSet;
-import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -53,7 +51,7 @@ import java.util.stream.Collectors;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
-import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
/** Test class for {@link org.apache.paimon.utils.ChainTableUtils}. */
public class ChainTableUtilsTest {
@@ -186,38 +184,6 @@ public class ChainTableUtilsTest {
Assertions.assertTrue(predicate.equals(expected));
}
- @Test
- public void testGeneratePartitionValues() {
- LinkedHashMap<String, String> partitionValues =
- ChainTableUtils.calPartValues(
- LocalDateTime.of(2023, 1, 1, 12, 0, 0),
- Arrays.asList("dt", "hour"),
- "$dt $hour:00:00",
- "yyyyMMdd HH:mm:ss");
- assertEquals(
- new LinkedHashMap<String, String>() {
- {
- put("dt", "20230101");
- put("hour", "12");
- }
- },
- partitionValues);
-
- partitionValues =
- ChainTableUtils.calPartValues(
- LocalDateTime.of(2023, 1, 1, 0, 0, 0),
- Arrays.asList("dt"),
- "$dt",
- "yyyyMMdd");
- assertEquals(
- new LinkedHashMap<String, String>() {
- {
- put("dt", "20230101");
- }
- },
- partitionValues);
- }
-
// ========================== Tests for
findFirstLatestPartitionsWithProjector
// ==========================
@@ -917,4 +883,116 @@ public class ChainTableUtilsTest {
private static Set<String> fileNames(ChainSplit split) {
return
split.dataFiles().stream().map(DataFileMeta::fileName).collect(Collectors.toSet());
}
+
+ @Test
+ public void testGetDeltaPartitionsWithHourMinuteGranularity() {
+ // partition keys: (region, dt, hour_minute), chain keys: (dt,
hour_minute)
+ RowType fullType =
+ RowType.builder()
+ .field("region", DataTypes.STRING().notNull())
+ .field("dt", DataTypes.STRING().notNull())
+ .field("hour_minute", DataTypes.STRING().notNull())
+ .build();
+
+ ChainPartitionProjector projector = new
ChainPartitionProjector(fullType, 2);
+
+ // Compare chain partition (dt, hour_minute) lexicographically
+ RecordComparator chainComparator = (a, b) ->
a.getString(1).compareTo(b.getString(1));
+
+ Options opts = new Options();
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_PATTERN, "$dtT$hour_minute");
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_FORMATTER, "yyyyMMdd'T'HHmm");
+ CoreOptions options = new CoreOptions(opts);
+
+ BinaryRow begin = row(Lists.newArrayList("CN", "20260609", "1010"));
+ BinaryRow end = row(Lists.newArrayList("CN", "20260609", "1015"));
+
+ List<BinaryRow> deltas =
+ ChainTableUtils.getDeltaPartitionsWithProjector(
+ begin, end, options, chainComparator, projector);
+
+ assertThat(deltas).hasSize(5);
+ for (BinaryRow delta : deltas) {
+ assertThat(getString(delta, 0)).isEqualTo("CN");
+ assertThat(getString(delta, 1)).isEqualTo("20260609");
+ }
+ assertThat(getString(deltas.get(0), 2)).isEqualTo("1011");
+ assertThat(getString(deltas.get(1), 2)).isEqualTo("1012");
+ assertThat(getString(deltas.get(2), 2)).isEqualTo("1013");
+ assertThat(getString(deltas.get(3), 2)).isEqualTo("1014");
+ assertThat(getString(deltas.get(4), 2)).isEqualTo("1015");
+ }
+
+ @Test
+ public void testGetDeltaPartitionsWithSeparateHourAndMinute() {
+ // partition keys: (region, dt, hour, minute), chain keys: (dt, hour,
minute)
+ RowType fullType =
+ RowType.builder()
+ .field("region", DataTypes.STRING().notNull())
+ .field("dt", DataTypes.STRING().notNull())
+ .field("hour", DataTypes.STRING().notNull())
+ .field("minute", DataTypes.STRING().notNull())
+ .build();
+
+ ChainPartitionProjector projector = new
ChainPartitionProjector(fullType, 3);
+
+ // Compare chain partition (dt, hour, minute) lexicographically
+ RecordComparator chainComparator = (a, b) ->
a.getString(2).compareTo(b.getString(2));
+
+ Options opts = new Options();
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_PATTERN,
"$dtT$hour$minute00");
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_FORMATTER,
"yyyyMMdd'T'HHmmss");
+ CoreOptions options = new CoreOptions(opts);
+
+ BinaryRow begin = row(Lists.newArrayList("CN", "20260609", "10",
"10"));
+ BinaryRow end = row(Lists.newArrayList("CN", "20260609", "10", "15"));
+
+ List<BinaryRow> deltas =
+ ChainTableUtils.getDeltaPartitionsWithProjector(
+ begin, end, options, chainComparator, projector);
+
+ assertThat(deltas).hasSize(5);
+ for (BinaryRow delta : deltas) {
+ assertThat(getString(delta, 0)).isEqualTo("CN");
+ assertThat(getString(delta, 1)).isEqualTo("20260609");
+ assertThat(getString(delta, 2)).isEqualTo("10");
+ }
+ assertThat(getString(deltas.get(0), 3)).isEqualTo("11");
+ assertThat(getString(deltas.get(1), 3)).isEqualTo("12");
+ assertThat(getString(deltas.get(2), 3)).isEqualTo("13");
+ assertThat(getString(deltas.get(3), 3)).isEqualTo("14");
+ assertThat(getString(deltas.get(4), 3)).isEqualTo("15");
+ }
+
+ @Test
+ public void testGetDeltaPartitionsExceedsMaxLimit() {
+ RowType partType = RowType.builder().field("dt",
DataTypes.STRING().notNull()).build();
+
+ Options opts = new Options();
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_PATTERN, "$dt");
+ opts.set(CoreOptions.PARTITION_TIMESTAMP_FORMATTER, "yyyyMMddHHmmss");
+ CoreOptions options = new CoreOptions(opts);
+
+ InternalRowPartitionComputer partitionComputer =
+ new InternalRowPartitionComputer(
+ options.partitionDefaultName(),
+ partType,
+ new String[] {"dt"},
+ options.legacyPartitionName());
+
+ BinaryRow begin = row(Lists.newArrayList("20250101000000"));
+ BinaryRow end = row(Lists.newArrayList("20260101000000"));
+
+ assertThatThrownBy(
+ () ->
+ ChainTableUtils.getDeltaPartitions(
+ begin,
+ end,
+ Collections.singletonList("dt"),
+ partType,
+ options,
+ partitionComputer))
+ .isInstanceOf(IllegalArgumentException.class)
+ .hasMessageContaining("Too many delta partitions generated");
+ }
}
diff --git
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkChainTableITCase.java
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkChainTableITCase.java
index f3a396c481..299868cef2 100644
---
a/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkChainTableITCase.java
+++
b/paimon-spark/paimon-spark-ut/src/test/java/org/apache/paimon/spark/SparkChainTableITCase.java
@@ -2422,6 +2422,63 @@ public class SparkChainTableITCase {
spark.close();
}
+ @Test
+ public void testChainTableWithMinuteLevelPartitions(@TempDir
java.nio.file.Path tempDir)
+ throws IOException {
+ Path warehousePath = new Path("file:" + tempDir.toString());
+ SparkSession.Builder builder =
createSparkSessionBuilder(warehousePath);
+ SparkSession spark = builder.getOrCreate();
+ spark.sql("CREATE DATABASE IF NOT EXISTS my_db1");
+ spark.sql("USE spark_catalog.my_db1");
+
+ spark.sql(
+ "CREATE TABLE `chain_test` (\n"
+ + " `t1` BIGINT COMMENT 't1',\n"
+ + " `t2` BIGINT COMMENT 't2',\n"
+ + " `t3` STRING COMMENT 't3'\n"
+ + ") PARTITIONED BY (`dt` STRING, `hr_min` STRING)\n"
+ + "TBLPROPERTIES (\n"
+ + " 'bucket-key' = 't1',\n"
+ + " 'primary-key' = 'dt,hr_min,t1',\n"
+ + " 'partition.timestamp-pattern' = '$dt
$hr_min:00',\n"
+ + " 'partition.timestamp-formatter' = 'yyyyMMdd
HH:mm:ss',\n"
+ + " 'chain-table.enabled' = 'true',\n"
+ + " 'bucket' = '1',\n"
+ + " 'merge-engine' = 'deduplicate',\n"
+ + " 'sequence.field' = 't2',\n"
+ + " 'chain-table.chain-partition-keys' =
'dt,hr_min'\n"
+ + ");");
+
+ setupChainTableBranches(spark, "chain_test");
+
+ spark.sql(
+ "INSERT INTO TABLE `chain_test$branch_snapshot` PARTITION (dt
= '20250810', hr_min='01:01') VALUES (3, 1, '3');");
+ spark.sql(
+ "INSERT INTO TABLE `chain_test$branch_snapshot` PARTITION (dt
= '20250810', hr_min='03:30') VALUES (4, 1, '4');");
+
+ spark.sql(
+ "INSERT INTO TABLE `chain_test$branch_delta` PARTITION (dt =
'20250810', hr_min='03:35') VALUES (5, 1, '5');");
+ spark.sql(
+ "INSERT INTO TABLE `chain_test$branch_delta` PARTITION (dt =
'20250810', hr_min='03:40') VALUES (6, 1, '6');");
+ spark.sql(
+ "INSERT INTO TABLE `chain_test$branch_delta` PARTITION (dt =
'20250810', hr_min='03:45') VALUES (7, 1, '7');");
+
+ assertThat(
+ spark
+ .sql(
+ "select * from `chain_test` where
dt='20250810' and hr_min='03:40'")
+ .collectAsList().stream()
+ .map(Row::toString)
+ .collect(Collectors.toList()))
+ .containsExactlyInAnyOrder(
+ "[4,1,4,20250810,03:40]",
+ "[5,1,5,20250810,03:40]",
+ "[6,1,6,20250810,03:40]");
+
+ spark.sql("DROP TABLE IF EXISTS `my_db1`.`chain_test`;");
+ spark.close();
+ }
+
@Test
public void testChainTableWithDeletionVectors(@TempDir java.nio.file.Path
tempDir)
throws IOException {