diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java index e83f5e81c1b7..4f1994f88d9a 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/AbstractFileStoreWrite.java @@ -37,7 +37,7 @@ import org.apache.paimon.memory.MemoryPoolFactory; import org.apache.paimon.metrics.MetricRegistry; import org.apache.paimon.operation.metrics.CompactionMetrics; -import org.apache.paimon.partition.PartitionTimeExtractor; +import org.apache.paimon.partition.PartitionTimeResolvable; import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.sink.CommitMessageImpl; import org.apache.paimon.types.RowType; @@ -714,15 +714,15 @@ public CompactionMetrics compactionMetrics() { private static class PartitionTimestampValidator { - private final PartitionTimeExtractor timeExtractor; + private final PartitionTimeResolvable timeResolver; private final RowDataToObjectArrayConverter partitionConverter; private final List partitionKeys; private PartitionTimestampValidator( - PartitionTimeExtractor timeExtractor, + PartitionTimeResolvable timeResolver, RowDataToObjectArrayConverter partitionConverter, List partitionKeys) { - this.timeExtractor = timeExtractor; + this.timeResolver = timeResolver; this.partitionConverter = partitionConverter; this.partitionKeys = partitionKeys; } @@ -738,7 +738,8 @@ private static PartitionTimestampValidator create( if ((timeFormatter != null || timePattern != null) && partitionType.getFieldCount() > 0) { return new PartitionTimestampValidator( - new PartitionTimeExtractor(timePattern, timeFormatter), + PartitionTimeResolvable.create( + partitionType.getFieldNames(), timePattern, timeFormatter), new RowDataToObjectArrayConverter(partitionType), partitionType.getFieldNames()); } @@ -748,7 +749,7 @@ private static PartitionTimestampValidator create( private void validate(BinaryRow partition) { Object[] array = partitionConverter.convert(partition); try { - timeExtractor.extract(partitionKeys, Arrays.asList(array)); + timeResolver.parsePartitionValues(Arrays.asList(array)); } catch (DateTimeParseException e) { String partitionInfo = IntStream.range(0, partitionKeys.size()) diff --git a/paimon-core/src/main/java/org/apache/paimon/operation/ChainTablePartitionExpire.java b/paimon-core/src/main/java/org/apache/paimon/operation/ChainTablePartitionExpire.java index 0e70854363ad..711ef326616e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/operation/ChainTablePartitionExpire.java +++ b/paimon-core/src/main/java/org/apache/paimon/operation/ChainTablePartitionExpire.java @@ -25,7 +25,7 @@ import org.apache.paimon.codegen.RecordComparator; import org.apache.paimon.data.BinaryRow; import org.apache.paimon.manifest.PartitionEntry; -import org.apache.paimon.partition.PartitionTimeExtractor; +import org.apache.paimon.partition.PartitionTimeResolver; import org.apache.paimon.table.FileStoreTable; import org.apache.paimon.table.PartitionModification; import org.apache.paimon.table.sink.BatchTableCommit; @@ -85,7 +85,7 @@ public class ChainTablePartitionExpire implements PartitionExpire { private final Duration checkInterval; private final FileStoreTable snapshotTable; private final FileStoreTable deltaTable; - private final PartitionTimeExtractor timeExtractor; + private final PartitionTimeResolver timeResolver; private final ChainPartitionProjector projector; private final RecordComparator chainPartitionComparator; private final InternalRowPartitionComputer partitionComputer; @@ -126,9 +126,11 @@ public ChainTablePartitionExpire( this.projector = new ChainPartitionProjector(partitionType, chainFieldCount); this.chainPartitionComparator = CodeGenUtils.newRecordComparator(projector.chainPartitionType().getFieldTypes()); - this.timeExtractor = - new PartitionTimeExtractor( - options.partitionTimestampPattern(), options.partitionTimestampFormatter()); + this.timeResolver = + new PartitionTimeResolver( + this.chainPartitionKeys, + options.partitionTimestampPattern(), + options.partitionTimestampFormatter()); this.partitionComputer = new InternalRowPartitionComputer( options.partitionDefaultName(), @@ -408,7 +410,7 @@ private LocalDateTime extractPartitionTime(BinaryRow partition) { for (String key : chainPartitionKeys) { chainValues.add(partValues.get(key)); } - return timeExtractor.extract(chainPartitionKeys, chainValues); + return timeResolver.parsePartitionValues(chainValues); } catch (Exception e) { LOG.warn("Failed to extract partition time from {}", partition, e); return null; diff --git a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeExtractor.java b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeExtractor.java deleted file mode 100644 index afc58c851317..000000000000 --- a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeExtractor.java +++ /dev/null @@ -1,134 +0,0 @@ -/* - * 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 javax.annotation.Nullable; - -import java.io.Serializable; -import java.time.LocalDate; -import java.time.LocalDateTime; -import java.time.LocalTime; -import java.time.format.DateTimeFormatter; -import java.time.format.DateTimeFormatterBuilder; -import java.time.format.DateTimeParseException; -import java.time.format.ResolverStyle; -import java.time.format.SignStyle; -import java.time.temporal.ChronoField; -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Locale; -import java.util.Objects; - -import static java.time.temporal.ChronoField.DAY_OF_MONTH; -import static java.time.temporal.ChronoField.HOUR_OF_DAY; -import static java.time.temporal.ChronoField.MINUTE_OF_HOUR; -import static java.time.temporal.ChronoField.MONTH_OF_YEAR; -import static java.time.temporal.ChronoField.SECOND_OF_MINUTE; -import static java.time.temporal.ChronoField.YEAR; - -/** Time extractor to extract time from partition values. */ -public class PartitionTimeExtractor implements Serializable { - - private static final long serialVersionUID = 1L; - - private static final DateTimeFormatter TIMESTAMP_FORMATTER = - new DateTimeFormatterBuilder() - .appendValue(YEAR, 1, 10, SignStyle.NORMAL) - .appendLiteral('-') - .appendValue(MONTH_OF_YEAR, 1, 2, SignStyle.NORMAL) - .appendLiteral('-') - .appendValue(DAY_OF_MONTH, 1, 2, SignStyle.NORMAL) - .optionalStart() - .appendLiteral(" ") - .appendValue(HOUR_OF_DAY, 1, 2, SignStyle.NORMAL) - .appendLiteral(':') - .appendValue(MINUTE_OF_HOUR, 1, 2, SignStyle.NORMAL) - .appendLiteral(':') - .appendValue(SECOND_OF_MINUTE, 1, 2, SignStyle.NORMAL) - .optionalStart() - .appendFraction(ChronoField.NANO_OF_SECOND, 1, 9, true) - .optionalEnd() - .optionalEnd() - .toFormatter() - .withResolverStyle(ResolverStyle.LENIENT); - - private static final DateTimeFormatter DATE_FORMATTER = - new DateTimeFormatterBuilder() - .appendValue(YEAR, 1, 10, SignStyle.NORMAL) - .appendLiteral('-') - .appendValue(MONTH_OF_YEAR, 1, 2, SignStyle.NORMAL) - .appendLiteral('-') - .appendValue(DAY_OF_MONTH, 1, 2, SignStyle.NORMAL) - .toFormatter() - .withResolverStyle(ResolverStyle.LENIENT); - - @Nullable private final String pattern; - @Nullable private final String formatter; - - public PartitionTimeExtractor(@Nullable String pattern, @Nullable String formatter) { - this.pattern = pattern; - this.formatter = formatter; - } - - public LocalDateTime extract(LinkedHashMap spec) { - return extract(new ArrayList<>(spec.keySet()), new ArrayList<>(spec.values())); - } - - public LocalDateTime extract(List partitionKeys, List partitionValues) { - String timestampString; - if (pattern == null) { - timestampString = partitionValues.get(0).toString(); - } else { - timestampString = pattern; - for (int i = 0; i < partitionKeys.size(); i++) { - timestampString = - timestampString.replaceAll( - "\\$" + partitionKeys.get(i), partitionValues.get(i).toString()); - } - } - return toLocalDateTime(timestampString, this.formatter); - } - - private static LocalDateTime toLocalDateTime( - String timestampString, @Nullable String formatterPattern) { - - if (formatterPattern == null) { - return PartitionTimeExtractor.toLocalDateTimeDefault(timestampString); - } - DateTimeFormatter dateTimeFormatter = - DateTimeFormatter.ofPattern(Objects.requireNonNull(formatterPattern), Locale.ROOT); - try { - return LocalDateTime.parse(timestampString, Objects.requireNonNull(dateTimeFormatter)); - } catch (DateTimeParseException e) { - return LocalDateTime.of( - LocalDate.parse(timestampString, Objects.requireNonNull(dateTimeFormatter)), - LocalTime.MIDNIGHT); - } - } - - public static LocalDateTime toLocalDateTimeDefault(String timestampString) { - try { - return LocalDateTime.parse(timestampString, TIMESTAMP_FORMATTER); - } catch (DateTimeParseException e) { - return LocalDateTime.of( - LocalDate.parse(timestampString, DATE_FORMATTER), LocalTime.MIDNIGHT); - } - } -} diff --git a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolvable.java b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolvable.java new file mode 100644 index 000000000000..7c4a246243de --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolvable.java @@ -0,0 +1,67 @@ +/* + * 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.io.Serializable; +import java.time.LocalDateTime; +import java.time.temporal.TemporalAmount; +import java.util.LinkedHashMap; +import java.util.List; + +/** + * Resolves partition values to/from timestamp and extracts the minimum time step. + * + *

Use {@link #create(List, String, String)} to obtain an instance. If both {@code pattern} and + * {@code formatter} are provided, a full pattern-based resolver is returned; otherwise a fallback + * resolver that handles the unconfigured case is returned. + */ +public interface PartitionTimeResolvable extends Serializable { + + /** Returns the partition keys used by this resolver. */ + List partitionKeys(); + + /** Parses partition column values into a {@link LocalDateTime}. */ + LocalDateTime parsePartitionValues(List partitionValues); + + /** Formats a {@link LocalDateTime} into partition column values. */ + default LinkedHashMap resolvePartitionValues(LocalDateTime dateTime) { + throw new UnsupportedOperationException( + "resolvePartitionValues is not supported by this resolver"); + } + + /** Extracts the minimum time step covered by the partition pattern and formatter. */ + default TemporalAmount extractMinStep() { + throw new UnsupportedOperationException("extractMinStep is not supported by this resolver"); + } + + /** + * Creates a {@link PartitionTimeResolvable}. + * + *

If both {@code pattern} and {@code formatter} are non-null, returns a normal {@link + * PartitionTimeResolver}. If either is null, returns a fallback resolver that handles the + * unconfigured case. + */ + static PartitionTimeResolvable create( + List partitionKeys, String pattern, String formatter) { + if (pattern == null || formatter == null) { + return PartitionTimeResolver.createFallback(partitionKeys, pattern, formatter); + } + return new PartitionTimeResolver(partitionKeys, pattern, formatter); + } +} 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 000000000000..e8548da71e76 --- /dev/null +++ b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionTimeResolver.java @@ -0,0 +1,593 @@ +/* + * 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 javax.annotation.Nullable; + +import java.text.ParsePosition; +import java.time.DateTimeException; +import java.time.Duration; +import java.time.LocalDate; +import java.time.LocalDateTime; +import java.time.LocalTime; +import java.time.Period; +import java.time.Year; +import java.time.YearMonth; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeFormatterBuilder; +import java.time.format.DateTimeParseException; +import java.time.format.ResolverStyle; +import java.time.format.SignStyle; +import java.time.temporal.ChronoField; +import java.time.temporal.TemporalAccessor; +import java.time.temporal.TemporalAmount; +import java.time.temporal.TemporalField; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Comparator; +import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Locale; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.stream.Collectors; + +import static org.apache.paimon.utils.Preconditions.checkArgument; + +/** + * Pattern-based implementation of {@link PartitionTimeResolvable}. It matches the user-provided + * timestamp pattern against the formatter and supports bidirectional conversion between partition + * values and {@link LocalDateTime}. + */ +public class PartitionTimeResolver implements PartitionTimeResolvable { + private static final Map FIELD_MAP = new HashMap<>(); + private final List partitionKeys; + private final String pattern; + private final String formatter; + private Map> patternFormatMappings; + private List patternTokens; + private List formatTokens; + + public PartitionTimeResolver(List partitionKeys, String pattern, String formatter) { + checkArgument(pattern != null, "pattern cannot be null"); + checkArgument(formatter != null, "formatter cannot be null"); + checkArgument(partitionKeys != null, "partitionKeys cannot be null"); + this.partitionKeys = partitionKeys; + this.pattern = pattern; + this.formatter = formatter; + init(); + } + + /** + * Creates a fallback resolver used when the user has not configured {@code + * partition.timestamp-pattern} or {@code partition.timestamp-formatter}. The fallback supports + * the common unconfigured case. + */ + static PartitionTimeResolvable createFallback( + List partitionKeys, String pattern, String formatter) { + return new FallbackPartitionTimeResolver(partitionKeys, pattern, formatter); + } + + static { + FIELD_MAP.put('y', ChronoField.YEAR); + FIELD_MAP.put('M', ChronoField.MONTH_OF_YEAR); + FIELD_MAP.put('d', ChronoField.DAY_OF_MONTH); + FIELD_MAP.put('H', ChronoField.HOUR_OF_DAY); + FIELD_MAP.put('h', ChronoField.CLOCK_HOUR_OF_AMPM); + 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(partitionKeys, pattern); + this.formatTokens = parseFormatter(); + boolean matched = matchRecursive(0, 0); + checkArgument( + matched, "Failed to match pattern '%s' to formatter '%s'", pattern, formatter); + } + + @Override + public List partitionKeys() { + return partitionKeys; + } + + /** + * 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 + */ + @Override + public TemporalAmount extractMinStep() { + List fieldTokens = + patternFormatMappings.values().stream() + .flatMap(Collection::stream) + .filter(token -> token instanceof TimeFieldToken) + .map(token -> (TimeFieldToken) token) + .collect(Collectors.toList()); + + Optional min = + fieldTokens.stream().min(Comparator.comparingInt(span -> span.field.ordinal())); + checkArgument(min.isPresent(), "No time unit found in variable ranges"); + ChronoField field = min.get().field; + return stepOf(field); + } + + /** + * Computes partition column values by formatting the given datetime and extracting each + * variable's segment according to the pattern-to-format mapping. + */ + @Override + public LinkedHashMap resolvePartitionValues(LocalDateTime dateTime) { + LinkedHashMap result = new LinkedHashMap<>(); + for (PatternToken patternToken : patternTokens) { + if (!patternToken.isVariable) { + continue; + } + String variableName = patternToken.token.substring(1); + List tokens = patternFormatMappings.get(patternToken); + int start = tokens.get(0).start; + int end = tokens.get(tokens.size() - 1).end; + DateTimeFormatter dateTimeFormatter = + DateTimeFormatter.ofPattern(formatter.substring(start, end), Locale.ROOT); + result.put(variableName, dateTime.format(dateTimeFormatter)); + } + return result; + } + + @Override + public LocalDateTime parsePartitionValues(List partitionValues) { + String timestampString = + buildTimestampString(partitionKeys, partitionValues, patternTokens); + DateTimeFormatter dateTimeFormatter = DateTimeFormatter.ofPattern(formatter, Locale.ROOT); + + Set fields = + formatTokens.stream() + .filter(t -> t instanceof TimeFieldToken) + .map(t -> ((TimeFieldToken) t).field) + .collect(Collectors.toSet()); + + if (fields.contains(ChronoField.HOUR_OF_DAY) + || fields.contains(ChronoField.CLOCK_HOUR_OF_AMPM) + || fields.contains(ChronoField.MINUTE_OF_HOUR) + || fields.contains(ChronoField.SECOND_OF_MINUTE)) { + return LocalDateTime.parse(timestampString, dateTimeFormatter); + } + if (fields.contains(ChronoField.DAY_OF_MONTH)) { + return LocalDate.parse(timestampString, dateTimeFormatter).atStartOfDay(); + } + if (fields.contains(ChronoField.MONTH_OF_YEAR)) { + return YearMonth.parse(timestampString, dateTimeFormatter).atDay(1).atStartOfDay(); + } + if (fields.contains(ChronoField.YEAR)) { + return Year.parse(timestampString, dateTimeFormatter) + .atMonth(1) + .atDay(1) + .atStartOfDay(); + } + throw new IllegalStateException("No time field found in formatter"); + } + + /** + * Builds the timestamp string by substituting partition column values into the pattern tokens. + */ + private static String buildTimestampString( + List partitionKeys, List partitionValues, List patternTokens) { + checkArgument(partitionValues != null, "Values cannot be null"); + + Map valueMap = new HashMap<>(); + for (int i = 0; i < partitionKeys.size(); i++) { + valueMap.put(partitionKeys.get(i), partitionValues.get(i)); + } + checkArgument(partitionValues.size() == valueMap.size(), "Values size mismatch"); + + StringBuilder timestampString = new StringBuilder(); + for (PatternToken token : patternTokens) { + if (token.isVariable) { + timestampString.append(valueMap.get(token.token.substring(1))); + } else { + timestampString.append(token.token); + } + } + return timestampString.toString(); + } + + /** Parses formatter into format tokens (time fields and literals). */ + private List parseFormatter() { + List 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++; + } + ChronoField field = FIELD_MAP.get(c); + tokens.add(new TimeFieldToken(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 static List parsePattern(List partitionKeys, String pattern) { + List sortedPartCols = + partitionKeys.stream() + .sorted(Comparator.reverseOrder()) + .collect(Collectors.toList()); + + List 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; + 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 (minFormattedLength(formatIdx, formatEndIdx) > patternToken.token.length()) { + continue; + } + 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; + 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 (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 : FIELD_MAP.values()) { + if (ta.isSupported(field)) { + try { + ta.get(field); + } catch (DateTimeException ignored) { + return false; + } + } + } + } catch (Exception ignored) { + return false; + } + return true; + } + + /** Minimum formatted length of the given format-token range; used to prune literal matches. */ + private int minFormattedLength(int startIdx, int endIdx) { + int length = 0; + for (int i = startIdx; i < endIdx; i++) { + FormatToken token = formatTokens.get(i); + if (token instanceof TimeFieldToken) { + TimeFieldToken t = (TimeFieldToken) token; + // Text month forms (MMMM/MMMMM) output variable-length strings whose length + // depends on the month value and locale. For example, MMMM outputs "May" (3) in + // English but "五月" (2) in Chinese; MMMMM can output a single character. + // Use a conservative lower bound of 1 to avoid pruning valid matches. + if (t.field == ChronoField.MONTH_OF_YEAR && token.getLength() >= 3) { + length += 1; + continue; + } + } + length += token.getLength(); + } + return length; + } + + private static TemporalAmount stepOf(ChronoField field) { + switch (field) { + case SECOND_OF_MINUTE: + return Duration.ofSeconds(1); + case MINUTE_OF_HOUR: + return Duration.ofMinutes(1); + case HOUR_OF_DAY: + case CLOCK_HOUR_OF_AMPM: + return Duration.ofHours(1); + case DAY_OF_MONTH: + return Duration.ofDays(1); + case MONTH_OF_YEAR: + return Period.ofMonths(1); + case YEAR: + return Period.ofYears(1); + default: + throw new IllegalStateException("Unsupported field: " + field); + } + } + + 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 TimeFieldToken extends FormatToken { + final ChronoField field; + + TimeFieldToken(ChronoField field, int start, int end) { + super(start, end); + this.field = field; + } + + @Override + public String toString() { + return String.format("TimeFieldToken{field=%s, start=%d, end=%d}", 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); + } + } + + /** + * Fallback resolver for the unconfigured case. When {@code partition.timestamp-pattern} or + * {@code partition.timestamp-formatter} is missing, pattern defaults to the first partition + * column and formatter defaults to {@code yyyy-MM-dd HH:mm:ss} / {@code yyyy-MM-dd}. + */ + private static class FallbackPartitionTimeResolver implements PartitionTimeResolvable { + private static final DateTimeFormatter TIMESTAMP_FORMATTER = + new DateTimeFormatterBuilder() + .appendValue(ChronoField.YEAR, 1, 10, SignStyle.NORMAL) + .appendLiteral('-') + .appendValue(ChronoField.MONTH_OF_YEAR, 1, 2, SignStyle.NORMAL) + .appendLiteral('-') + .appendValue(ChronoField.DAY_OF_MONTH, 1, 2, SignStyle.NORMAL) + .optionalStart() + .appendLiteral(" ") + .appendValue(ChronoField.HOUR_OF_DAY, 1, 2, SignStyle.NORMAL) + .appendLiteral(':') + .appendValue(ChronoField.MINUTE_OF_HOUR, 1, 2, SignStyle.NORMAL) + .appendLiteral(':') + .appendValue(ChronoField.SECOND_OF_MINUTE, 1, 2, SignStyle.NORMAL) + .optionalStart() + .appendFraction(ChronoField.NANO_OF_SECOND, 1, 9, true) + .optionalEnd() + .optionalEnd() + .toFormatter() + .withResolverStyle(ResolverStyle.LENIENT); + + private static final DateTimeFormatter DATE_FORMATTER = + new DateTimeFormatterBuilder() + .appendValue(ChronoField.YEAR, 1, 10, SignStyle.NORMAL) + .appendLiteral('-') + .appendValue(ChronoField.MONTH_OF_YEAR, 1, 2, SignStyle.NORMAL) + .appendLiteral('-') + .appendValue(ChronoField.DAY_OF_MONTH, 1, 2, SignStyle.NORMAL) + .toFormatter() + .withResolverStyle(ResolverStyle.LENIENT); + + private final List partitionKeys; + @Nullable private final String pattern; + @Nullable private final String formatter; + + FallbackPartitionTimeResolver( + List partitionKeys, @Nullable String pattern, @Nullable String formatter) { + checkArgument(partitionKeys != null, "partitionKeys cannot be null"); + this.partitionKeys = partitionKeys; + this.pattern = pattern; + this.formatter = formatter; + } + + @Override + public List partitionKeys() { + return partitionKeys; + } + + @Override + public LocalDateTime parsePartitionValues(List partitionValues) { + checkArgument(partitionValues != null, "Values cannot be null"); + String timestampString; + if (pattern == null) { + timestampString = partitionValues.get(0).toString(); + } else { + timestampString = + buildTimestampString( + partitionKeys, + partitionValues, + parsePattern(partitionKeys, pattern)); + } + return toLocalDateTime(timestampString, formatter); + } + + private static LocalDateTime toLocalDateTime( + String timestampString, @Nullable String formatterPattern) { + if (formatterPattern == null) { + try { + return LocalDateTime.parse(timestampString, TIMESTAMP_FORMATTER); + } catch (DateTimeParseException e) { + return LocalDateTime.of( + LocalDate.parse(timestampString, DATE_FORMATTER), LocalTime.MIDNIGHT); + } + } + DateTimeFormatter dateTimeFormatter = + DateTimeFormatter.ofPattern(formatterPattern, Locale.ROOT); + try { + return LocalDateTime.parse(timestampString, dateTimeFormatter); + } catch (DateTimeParseException e) { + return LocalDateTime.of( + LocalDate.parse(timestampString, dateTimeFormatter), LocalTime.MIDNIGHT); + } + } + } +} diff --git a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionValuesTimeExpireStrategy.java b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionValuesTimeExpireStrategy.java index 6238817da4e1..7c2b69d20b3e 100644 --- a/paimon-core/src/main/java/org/apache/paimon/partition/PartitionValuesTimeExpireStrategy.java +++ b/paimon-core/src/main/java/org/apache/paimon/partition/PartitionValuesTimeExpireStrategy.java @@ -50,13 +50,14 @@ public class PartitionValuesTimeExpireStrategy extends PartitionExpireStrategy { private static final Logger LOG = LoggerFactory.getLogger(PartitionValuesTimeExpireStrategy.class); - private final PartitionTimeExtractor timeExtractor; + private final PartitionTimeResolvable timeResolver; public PartitionValuesTimeExpireStrategy(CoreOptions options, RowType partitionType) { super(partitionType, options.partitionDefaultName()); String timePattern = options.partitionTimestampPattern(); String timeFormatter = options.partitionTimestampFormatter(); - this.timeExtractor = new PartitionTimeExtractor(timePattern, timeFormatter); + this.timeResolver = + PartitionTimeResolvable.create(partitionKeys, timePattern, timeFormatter); } @Override @@ -83,7 +84,7 @@ private PartitionValuesTimePredicate(LocalDateTime expireDateTime) { public boolean test(BinaryRow partition) { Object[] array = convertPartition(partition); try { - LocalDateTime partTime = timeExtractor.extract(partitionKeys, Arrays.asList(array)); + LocalDateTime partTime = timeResolver.parsePartitionValues(Arrays.asList(array)); return expireDateTime.isAfter(partTime); } catch (DateTimeParseException e) { LOG.warn( 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 4a7bb1b073a5..d3d23d77cba2 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.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; @@ -40,7 +40,7 @@ import javax.annotation.Nullable; import java.time.LocalDateTime; -import java.time.format.DateTimeFormatter; +import java.time.temporal.TemporalAmount; import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; @@ -52,8 +52,6 @@ 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; @@ -92,61 +90,33 @@ public static List getDeltaPartitions( List partitionColumns, RowType partType, CoreOptions options, - RecordComparator partitionComparator, InternalRowPartitionComputer partitionComputer) { InternalRowSerializer serializer = new InternalRowSerializer(partType); List deltaPartitions = new ArrayList<>(); - boolean isDailyPartition = partitionColumns.size() == 1; List startPartitionValues = new ArrayList<>(partitionComputer.generatePartValues(beginPartition).values()); List 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 stratPartitionTime = timeResolver.parsePartitionValues(startPartitionValues); + LocalDateTime endPartitionTime = timeResolver.parsePartitionValues(endPartitionValues); + TemporalAmount step = timeResolver.extractMinStep(); + LocalDateTime candidateTime = stratPartitionTime.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 +153,6 @@ public static Predicate createLinearPredicate( return PredicateBuilder.and(fieldPredicates); } - public static LinkedHashMap calPartValues( - LocalDateTime dateTime, - List 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 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 keyValues = new HashMap<>(); - for (int i = 0; i < keyOrder.size(); i++) { - keyValues.put(keyOrder.get(i), valueMatcher.group(i + 1)); - } - List values = - partitionKeys.stream() - .map(key -> keyValues.getOrDefault(key, "")) - .collect(Collectors.toList()); - LinkedHashMap 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 +257,6 @@ public static List getDeltaPartitionsWithProjector( 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/operation/ChainTablePartitionExpireTest.java b/paimon-core/src/test/java/org/apache/paimon/operation/ChainTablePartitionExpireTest.java index c8440a3506f7..2b39a8f0cc3b 100644 --- a/paimon-core/src/test/java/org/apache/paimon/operation/ChainTablePartitionExpireTest.java +++ b/paimon-core/src/test/java/org/apache/paimon/operation/ChainTablePartitionExpireTest.java @@ -41,6 +41,7 @@ import org.apache.paimon.table.sink.CommitMessage; import org.apache.paimon.table.sink.StreamTableWrite; import org.apache.paimon.table.sink.TableCommitImpl; +import org.apache.paimon.types.DataField; import org.apache.paimon.types.DataTypes; import org.apache.paimon.types.RowType; @@ -48,6 +49,8 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; +import javax.annotation.Nullable; + import java.time.Duration; import java.time.LocalDateTime; import java.util.ArrayList; @@ -57,6 +60,7 @@ import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.function.Function; import java.util.stream.Collectors; import static org.assertj.core.api.Assertions.assertThat; @@ -115,8 +119,6 @@ public void testExpireWithSinglePartitionKey() throws Exception { write(deltaTable, "20250215", "v7"); write(deltaTable, "20250315", "v8"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); assertThat(listPartitions(snapshotTable)) .containsExactlyInAnyOrder("20250101", "20250201", "20250301"); assertThat(listPartitions(deltaTable)) @@ -155,9 +157,6 @@ public void testNoExpireWhenOnlyOneSnapshotBeforeCutoff() throws Exception { write(snapshotTable, "20250315", "v2"); write(deltaTable, "20250205", "v3"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - ChainTablePartitionExpire expire = newChainExpire(snapshotTable, deltaTable, Duration.ofDays(30), false); expire.setLastCheck(LocalDateTime.of(2025, 1, 1, 0, 0)); @@ -189,9 +188,6 @@ public void testExpireMultipleSegments() throws Exception { write(deltaTable, "20250120", "v6"); write(deltaTable, "20250210", "v7"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 // Snapshots before cutoff: 20250101, 20250115, 20250201 (3 snapshots) // Anchor = 20250201 (kept), expire S(20250101), S(20250115) @@ -244,9 +240,6 @@ public void testCheckIntervalPreventsExpire() throws Exception { write(snapshotTable, "20250101", "v1"); write(snapshotTable, "20250201", "v2"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - ChainTablePartitionExpire expire = new ChainTablePartitionExpire( Duration.ofDays(30), @@ -288,9 +281,6 @@ public void testMaxExpireNumLimitsSegments() throws Exception { write(deltaTable, "20250120", "v6"); write(deltaTable, "20250210", "v7"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // maxExpireNum=1 means only 1 segment: Segment1={S(0101), d(0105)} ChainTablePartitionExpire expire = new ChainTablePartitionExpire( @@ -343,9 +333,6 @@ public void testExpireWithGroupPartition() throws Exception { writeGrouped(snapshotTable, "EU", "20250320", "v5"); writeGrouped(deltaTable, "EU", "20250220", "d3"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 // Group "US": snapshots before cutoff = [0101, 0201]. Anchor = 0201 (kept). // Expire: S(0101), delta(0110) (before anchor 0201) @@ -363,8 +350,10 @@ public void testExpireWithGroupPartition() throws Exception { snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); deltaTable = loadTable(tablePath).switchToBranch("delta"); - List snapshotParts = listGroupedPartitions(snapshotTable); - List deltaParts = listGroupedPartitions(deltaTable); + List snapshotParts = + listPartitions(snapshotTable, p -> p.getString(0) + "|" + p.getString(1)); + List deltaParts = + listPartitions(deltaTable, p -> p.getString(0) + "|" + p.getString(1)); assertThat(snapshotParts).contains("US|20250201", "US|20250301"); assertThat(snapshotParts).doesNotContain("US|20250101"); @@ -389,9 +378,6 @@ public void testUsesBranchSpecificPartitionModifications() throws Exception { write(deltaTable, "20250110", "d1"); write(deltaTable, "20250210", "d2"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - RecordingPartitionModification snapshotModification = new RecordingPartitionModification(); RecordingPartitionModification deltaModification = new RecordingPartitionModification(); @@ -432,9 +418,6 @@ public void testIsValueAllExpiredReturnsFalseForAnchor() throws Exception { write(snapshotTable, "20250201", "v2"); write(snapshotTable, "20250301", "v3"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // expirationTime = 30d, "now" = 2025-03-31 → cutoff = 2025-03-01 // Snapshots before cutoff: S(0101), S(0201). Anchor = S(0201) (kept). ChainTablePartitionExpire expire = @@ -466,9 +449,6 @@ public void testIsValueAllExpiredReturnsFalseWhenTooFewSnapshots() throws Except write(snapshotTable, "20250201", "v1"); write(snapshotTable, "20250315", "v2"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 // Only S(0201) before cutoff → < 2, nothing can expire ChainTablePartitionExpire expire = @@ -491,9 +471,6 @@ public void testIsValueAllExpiredReturnsFalseForDeltaOnlyGroup() throws Exceptio write(deltaTable, "20250101", "d1"); write(deltaTable, "20250201", "d2"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 // No snapshot boundary exists, so delta-only partitions are retained. ChainTablePartitionExpire expire = @@ -522,9 +499,6 @@ public void testIsValueAllExpiredWithGroupPartitions() throws Exception { writeGrouped(snapshotTable, "EU", "20250215", "v4"); writeGrouped(snapshotTable, "EU", "20250320", "v5"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 ChainTablePartitionExpire expire = newChainExpire(snapshotTable, deltaTable, Duration.ofDays(30), true); @@ -558,9 +532,6 @@ public void testIsValueAllExpiredReturnsFalseForPartitionsAfterCutoff() throws E write(snapshotTable, "20250101", "v1"); write(snapshotTable, "20250315", "v2"); - snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); - deltaTable = loadTable(tablePath).switchToBranch("delta"); - // cutoff = 2025-03-31 - 30d = 2025-03-01 // S(0315) is after cutoff → not expired at all ChainTablePartitionExpire expire = @@ -571,6 +542,94 @@ public void testIsValueAllExpiredReturnsFalseForPartitionsAfterCutoff() throws E assertThat(expire.isValueAllExpired(Collections.singletonList(afterCutoff), now)).isFalse(); } + @Test + public void testExpireWithMinuteGranularity() throws Exception { + Path tablePath = tablePath("minute_granularity_expire"); + String pattern = "$y$m$dT$h$min00"; + String formatter = "yyyyMMdd'T'HHmmss"; + Map chainOptions = + buildOptions(pattern, formatter, "15 min", "y,m,d,h,min"); + createChainTable( + tablePath, + RowType.of( + new org.apache.paimon.types.DataType[] { + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING() + }, + new String[] {"y", "m", "d", "h", "min", "pk", "v"}) + .getFields(), + Arrays.asList("y", "m", "d", "h", "min"), + Arrays.asList("pk", "y", "m", "d", "h", "min"), + chainOptions); + + FileStoreTable mainTable = loadTable(tablePath); + FileStoreTable snapshotTable = mainTable.switchToBranch("snapshot"); + FileStoreTable deltaTable = mainTable.switchToBranch("delta"); + + // Snapshot partitions every hour. + write(snapshotTable, "2025", "01", "01", "10", "00", "v1", "v1"); + write(snapshotTable, "2025", "01", "01", "11", "00", "v2", "v2"); + write(snapshotTable, "2025", "01", "01", "12", "00", "v3", "v3"); + + // Delta partitions every 15 minutes. + write(deltaTable, "2025", "01", "01", "10", "15", "d1", "d1"); + write(deltaTable, "2025", "01", "01", "10", "30", "d2", "d2"); + write(deltaTable, "2025", "01", "01", "11", "15", "d3", "d3"); + write(deltaTable, "2025", "01", "01", "11", "30", "d4", "d4"); + write(deltaTable, "2025", "01", "01", "12", "15", "d5", "d5"); + + snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); + deltaTable = loadTable(tablePath).switchToBranch("delta"); + + // cutoff = 2025-01-01 12:15 - 15min = 2025-01-01 12:00 + // Snapshots before cutoff: 10:00, 11:00. Anchor = 11:00 (kept). + // Expire segment: S(10:00) + deltas in [10:00, 11:00): d(10:15), d(10:30). + ChainTablePartitionExpire expire = newChainExpire(snapshotTable, deltaTable, chainOptions); + expire.setLastCheck(LocalDateTime.of(2025, 1, 1, 0, 0)); + List> expired = + expire.expire(LocalDateTime.of(2025, 1, 1, 12, 15), Long.MAX_VALUE); + + assertThat(expired).isNotNull(); + assertThat(expired).hasSize(3); + + snapshotTable = loadTable(tablePath).switchToBranch("snapshot"); + deltaTable = loadTable(tablePath).switchToBranch("delta"); + assertThat( + listPartitions( + snapshotTable, + p -> + p.getString(0) + + "," + + p.getString(1) + + "," + + p.getString(2) + + "," + + p.getString(3) + + "," + + p.getString(4))) + .containsExactlyInAnyOrder("2025,01,01,11,00", "2025,01,01,12,00"); + assertThat( + listPartitions( + deltaTable, + p -> + p.getString(0) + + "," + + p.getString(1) + + "," + + p.getString(2) + + "," + + p.getString(3) + + "," + + p.getString(4))) + .containsExactlyInAnyOrder( + "2025,01,01,11,15", "2025,01,01,11,30", "2025,01,01,12,15"); + } + // ========== Helper methods ========== private BinaryRow findPartition(FileStoreTable table, String dtValue) { @@ -636,11 +695,13 @@ public void testRollbackToAsLatestRejectedWhenDroppingAnchorSnapshotPartition() .isEqualTo(2L); } - private Path tablePath(String tableName) { - return new Path(tempDir.toUri().toString(), tableName); - } - - private void createChainTable(Path tablePath, boolean withGroupPartition) throws Exception { + private void createChainTable( + Path tablePath, + List fields, + List partitionKeys, + List primaryKeys, + Map chainOptions) + throws Exception { LocalFileIO fileIO = LocalFileIO.create(); SchemaManager schemaManager = new SchemaManager(fileIO, tablePath); @@ -649,59 +710,64 @@ private void createChainTable(Path tablePath, boolean withGroupPartition) throws options.put("merge-engine", "deduplicate"); options.put("sequence.field", "v"); - Schema schema; - if (withGroupPartition) { - schema = - new Schema( - RowType.of( - new org.apache.paimon.types.DataType[] { - DataTypes.STRING(), - DataTypes.STRING(), - DataTypes.STRING(), - DataTypes.STRING() - }, - new String[] {"region", "dt", "pk", "v"}) - .getFields(), - Arrays.asList("region", "dt"), - Arrays.asList("pk", "region", "dt"), - options, - ""); - } else { - schema = - new Schema( - RowType.of( - new org.apache.paimon.types.DataType[] { - DataTypes.STRING(), - DataTypes.STRING(), - DataTypes.STRING() - }, - new String[] {"dt", "pk", "v"}) - .getFields(), - Collections.singletonList("dt"), - Arrays.asList("pk", "dt"), - options, - ""); - } + Schema schema = new Schema(fields, partitionKeys, primaryKeys, options, ""); schemaManager.createTable(schema); FileStoreTable mainTable = loadTable(tablePath); mainTable.createBranch("snapshot"); mainTable.createBranch("delta"); - List chainTableOptions = - Arrays.asList( - SchemaChange.setOption("chain-table.enabled", "true"), - SchemaChange.setOption("scan.fallback-snapshot-branch", "snapshot"), - SchemaChange.setOption("scan.fallback-delta-branch", "delta"), - SchemaChange.setOption("partition.timestamp-pattern", "$dt"), - SchemaChange.setOption("partition.timestamp-formatter", "yyyyMMdd")); + List changes = + new ArrayList<>( + Arrays.asList( + SchemaChange.setOption("chain-table.enabled", "true"), + SchemaChange.setOption("scan.fallback-snapshot-branch", "snapshot"), + SchemaChange.setOption("scan.fallback-delta-branch", "delta"))); + for (Map.Entry entry : chainOptions.entrySet()) { + changes.add(SchemaChange.setOption(entry.getKey(), entry.getValue())); + } + schemaManager.commitChanges(changes); + new SchemaManager(fileIO, tablePath, "snapshot").commitChanges(changes); + new SchemaManager(fileIO, tablePath, "delta").commitChanges(changes); + } + + private Path tablePath(String tableName) { + return new Path(tempDir.toUri().toString(), tableName); + } + + private void createChainTable(Path tablePath, boolean withGroupPartition) throws Exception { + Map chainOptions = new HashMap<>(); + chainOptions.put("partition.timestamp-pattern", "$dt"); + chainOptions.put("partition.timestamp-formatter", "yyyyMMdd"); if (withGroupPartition) { - chainTableOptions = new java.util.ArrayList<>(chainTableOptions); - chainTableOptions.add(SchemaChange.setOption("chain-table.chain-partition-keys", "dt")); + chainOptions.put("chain-table.chain-partition-keys", "dt"); + createChainTable( + tablePath, + RowType.of( + new org.apache.paimon.types.DataType[] { + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING(), + DataTypes.STRING() + }, + new String[] {"region", "dt", "pk", "v"}) + .getFields(), + Arrays.asList("region", "dt"), + Arrays.asList("pk", "region", "dt"), + chainOptions); + } else { + createChainTable( + tablePath, + RowType.of( + new org.apache.paimon.types.DataType[] { + DataTypes.STRING(), DataTypes.STRING(), DataTypes.STRING() + }, + new String[] {"dt", "pk", "v"}) + .getFields(), + Collections.singletonList("dt"), + Arrays.asList("pk", "dt"), + chainOptions); } - schemaManager.commitChanges(chainTableOptions); - new SchemaManager(fileIO, tablePath, "snapshot").commitChanges(chainTableOptions); - new SchemaManager(fileIO, tablePath, "delta").commitChanges(chainTableOptions); } private FileStoreTable loadTable(Path tablePath) { @@ -715,32 +781,24 @@ private FileStoreTable loadTable(Path tablePath) { } private void write(FileStoreTable table, String dt, String v) throws Exception { - StreamTableWrite write = - table.copy(Collections.singletonMap(CoreOptions.WRITE_ONLY.key(), "true")) - .newWrite(commitUser); - write.write( - GenericRow.of( - BinaryString.fromString(dt), - BinaryString.fromString(v), - BinaryString.fromString(v))); - TableCommitImpl commit = table.newCommit(commitUser); - List commitMessages = write.prepareCommit(true, 0); - commit.commit(0, commitMessages); - write.close(); - commit.close(); + write(table, dt, v, v); } private void writeGrouped(FileStoreTable table, String region, String dt, String v) throws Exception { + write(table, region, dt, v, v); + } + + private void write(FileStoreTable table, Object... values) throws Exception { + for (int i = 0; i < values.length; i++) { + if (values[i] != null && values[i] instanceof String) { + values[i] = BinaryString.fromString((String) values[i]); + } + } StreamTableWrite write = table.copy(Collections.singletonMap(CoreOptions.WRITE_ONLY.key(), "true")) .newWrite(commitUser); - write.write( - GenericRow.of( - BinaryString.fromString(region), - BinaryString.fromString(dt), - BinaryString.fromString(v), - BinaryString.fromString(v))); + write.write(GenericRow.of(values)); TableCommitImpl commit = table.newCommit(commitUser); List commitMessages = write.prepareCommit(true, 0); commit.commit(0, commitMessages); @@ -749,30 +807,39 @@ private void writeGrouped(FileStoreTable table, String region, String dt, String } private List listPartitions(FileStoreTable table) { - return table.newSnapshotReader().partitionEntries().stream() - .map(PartitionEntry::partition) - .map(p -> p.getString(0).toString()) - .sorted() - .collect(Collectors.toList()); + return listPartitions(table, p -> p.getString(0).toString()); } - private List listGroupedPartitions(FileStoreTable table) { + private List listPartitions( + FileStoreTable table, Function partitionToString) { return table.newSnapshotReader().partitionEntries().stream() .map(PartitionEntry::partition) - .map(p -> p.getString(0).toString() + "|" + p.getString(1).toString()) + .map(partitionToString) .sorted() .collect(Collectors.toList()); } private Map buildOptions(Duration expirationTime, boolean withGroupPartition) { + return buildOptions( + "$dt", + "yyyyMMdd", + expirationTime.toDays() + " d", + withGroupPartition ? "dt" : null); + } + + private Map buildOptions( + String timestampPattern, + String timestampFormatter, + String expirationTime, + @Nullable String chainPartitionKeys) { Map opts = new HashMap<>(); - opts.put("partition.timestamp-pattern", "$dt"); - opts.put("partition.timestamp-formatter", "yyyyMMdd"); + opts.put("partition.timestamp-pattern", timestampPattern); + opts.put("partition.timestamp-formatter", timestampFormatter); opts.put("scan.fallback-snapshot-branch", "snapshot"); opts.put("scan.fallback-delta-branch", "delta"); - opts.put(CoreOptions.PARTITION_EXPIRATION_TIME.key(), expirationTime.toDays() + " d"); - if (withGroupPartition) { - opts.put("chain-table.chain-partition-keys", "dt"); + opts.put(CoreOptions.PARTITION_EXPIRATION_TIME.key(), expirationTime); + if (chainPartitionKeys != null) { + opts.put("chain-table.chain-partition-keys", chainPartitionKeys); } return opts; } @@ -782,12 +849,18 @@ private ChainTablePartitionExpire newChainExpire( FileStoreTable deltaTable, Duration expirationTime, boolean withGroupPartition) { + return newChainExpire( + snapshotTable, deltaTable, buildOptions(expirationTime, withGroupPartition)); + } + + private ChainTablePartitionExpire newChainExpire( + FileStoreTable snapshotTable, FileStoreTable deltaTable, Map options) { return new ChainTablePartitionExpire( - expirationTime, + CoreOptions.fromMap(options).partitionExpireTime(), Duration.ZERO, snapshotTable, deltaTable, - CoreOptions.fromMap(buildOptions(expirationTime, withGroupPartition)), + CoreOptions.fromMap(options), snapshotTable.schema().logicalPartitionType(), false, Integer.MAX_VALUE, diff --git a/paimon-core/src/test/java/org/apache/paimon/partition/FallbackPartitionTimeResolverTest.java b/paimon-core/src/test/java/org/apache/paimon/partition/FallbackPartitionTimeResolverTest.java new file mode 100644 index 000000000000..3b4b93fc98b8 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/partition/FallbackPartitionTimeResolverTest.java @@ -0,0 +1,106 @@ +/* + * 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.testutils.assertj.PaimonAssertions; + +import org.junit.jupiter.api.Test; + +import java.time.LocalDateTime; +import java.time.format.DateTimeParseException; +import java.util.Arrays; +import java.util.Collections; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Test for {@link PartitionTimeResolvable} fallback behavior (unconfigured pattern/formatter). */ +public class FallbackPartitionTimeResolverTest { + + @Test + public void testDefault() { + PartitionTimeResolvable resolver = + PartitionTimeResolvable.create(Collections.emptyList(), null, null); + assertThat(resolver.parsePartitionValues(Collections.singletonList("2023-01-01 20:08:08"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T20:08:08")); + + assertThat( + resolver.parsePartitionValues( + Collections.singletonList("2023-01-01 20:08:08.12"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T20:08:08.12")); + + assertThat(resolver.parsePartitionValues(Collections.singletonList("2023-1-1 20:08:08"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T20:08:08")); + + assertThat(resolver.parsePartitionValues(Collections.singletonList("2023-01-01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + assertThat(resolver.parsePartitionValues(Collections.singletonList("2023-1-1"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + } + + @Test + public void testPattern() { + PartitionTimeResolvable resolver = + PartitionTimeResolvable.create( + Arrays.asList("year", "month", "day"), + "$year-$month-$day 00:00:00.12", + null); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023", "01", "01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00.12")); + + resolver = + PartitionTimeResolvable.create( + Arrays.asList("year", "month", "day", "hour"), + "$year-$month-$day $hour:00:00", + null); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023", "01", "01", "01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T01:00:00")); + + resolver = PartitionTimeResolvable.create(Arrays.asList("other", "dt"), "$dt", null); + assertThat(resolver.parsePartitionValues(Arrays.asList("dummy", "2023-01-01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + // One column name is a prefix of another ("t" vs "t1"). The longest match must be used + // so that "$t1" is not broken by "$t" replacement. + resolver = PartitionTimeResolvable.create(Arrays.asList("t", "t1"), "$t $t1", null); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-01-01", "00:00:00"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + } + + @Test + public void testFormatter() { + PartitionTimeResolvable resolver = + PartitionTimeResolvable.create(Collections.emptyList(), null, "yyyyMMdd"); + assertThat(resolver.parsePartitionValues(Collections.singletonList("20230101"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + } + + @Test + public void testExtractNonDateFormattedPartition() { + PartitionTimeResolvable resolver = + PartitionTimeResolvable.create(Collections.singletonList("ds"), "$ds", "yyyyMMdd"); + assertThatThrownBy( + () -> resolver.parsePartitionValues(Collections.singletonList("unknown"))) + .satisfies( + PaimonAssertions.anyCauseMatches( + DateTimeParseException.class, + "Text 'unknown' could not be parsed at index 0")); + } +} diff --git a/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeExtractorTest.java b/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeExtractorTest.java deleted file mode 100644 index 3f6cff6ee55d..000000000000 --- a/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeExtractorTest.java +++ /dev/null @@ -1,108 +0,0 @@ -/* - * 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.testutils.assertj.PaimonAssertions; - -import org.junit.jupiter.api.Test; - -import java.time.LocalDateTime; -import java.time.format.DateTimeParseException; -import java.util.Arrays; -import java.util.Collections; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatThrownBy; - -/** Test for {@link PartitionTimeExtractor}. */ -public class PartitionTimeExtractorTest { - - @Test - public void testDefault() { - PartitionTimeExtractor extractor = new PartitionTimeExtractor(null, null); - assertThat( - extractor.extract( - Collections.emptyList(), - Collections.singletonList("2023-01-01 20:08:08"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T20:08:08")); - - assertThat( - extractor.extract( - Collections.emptyList(), - Collections.singletonList("2023-1-1 20:08:08"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T20:08:08")); - - assertThat( - extractor.extract( - Collections.emptyList(), Collections.singletonList("2023-01-01"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); - - assertThat( - extractor.extract( - Collections.emptyList(), Collections.singletonList("2023-1-1"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); - } - - @Test - public void testPattern() { - PartitionTimeExtractor extractor = - new PartitionTimeExtractor("$year-$month-$day 00:00:00", null); - assertThat( - extractor.extract( - Arrays.asList("year", "month", "day"), - Arrays.asList("2023", "01", "01"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); - - extractor = new PartitionTimeExtractor("$year-$month-$day $hour:00:00", null); - assertThat( - extractor.extract( - Arrays.asList("year", "month", "day", "hour"), - Arrays.asList("2023", "01", "01", "01"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T01:00:00")); - - extractor = new PartitionTimeExtractor("$dt", null); - assertThat( - extractor.extract( - Arrays.asList("other", "dt"), Arrays.asList("dummy", "2023-01-01"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); - } - - @Test - public void testFormatter() { - PartitionTimeExtractor extractor = new PartitionTimeExtractor(null, "yyyyMMdd"); - assertThat( - extractor.extract( - Collections.emptyList(), Collections.singletonList("20230101"))) - .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); - } - - @Test - public void testExtractNonDateFormattedPartition() { - PartitionTimeExtractor extractor = new PartitionTimeExtractor("$ds", "yyyyMMdd"); - assertThatThrownBy( - () -> - extractor.extract( - Collections.singletonList("ds"), - Collections.singletonList("unknown"))) - .satisfies( - PaimonAssertions.anyCauseMatches( - DateTimeParseException.class, - "Text 'unknown' could not be parsed at index 0")); - } -} 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 000000000000..ce831e77f208 --- /dev/null +++ b/paimon-core/src/test/java/org/apache/paimon/partition/PartitionTimeResolverTest.java @@ -0,0 +1,479 @@ +/* + * 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.testutils.assertj.PaimonAssertions; +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 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 testExtractMinStepWithDuration() { + 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)); + assertThat(extractMinStep("12$b", "HHmmss", "b")).isEqualTo(Duration.ofSeconds(1)); + + assertThat(extractMinStep("$a", "yyyyMMddhh", "a")).isEqualTo(Duration.ofHours(1)); + assertThat(extractMinStep("$date $time", "yyyyMMdd hh", "date", "time")) + .isEqualTo(Duration.ofHours(1)); + assertThat(extractMinStep("$a $b", "yyyyMMdd hh", "a", "b")).isEqualTo(Duration.ofHours(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)); + + assertThatThrownBy(() -> extractMinStep("$dt", "yyyyMMddHHmmssSSS", "dt")) + .satisfies( + PaimonAssertions.anyCauseMatches( + IllegalArgumentException.class, + "Unsupported formatter pattern letter 'S' in formatter: yyyyMMddHHmmssSSS.")); + } + + @Test + public void testExtractMinStepWithPeriod() { + 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("$aJan", "yyMMM", "a")).isEqualTo(Period.ofYears(1)); + assertThat(extractMinStep("$a1", "yyM", "a")).isEqualTo(Period.ofYears(1)); + assertThat(extractMinStep("$a12", "yyM", "a")).isEqualTo(Period.ofYears(1)); + } + + @Test + public void testResolvePartitionValues() { + Map 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); + + partitionValues = + new PartitionTimeResolver(Arrays.asList("y", "m"), "$y$m", "yyMMM") + .resolvePartitionValues(LocalDateTime.of(2023, 12, 1, 11, 2, 3)); + assertEquals(ImmutableMap.of("y", "23", "m", "Dec"), partitionValues); + + partitionValues = + new PartitionTimeResolver(Arrays.asList("ym"), "$ym", "yy-M") + .resolvePartitionValues(LocalDateTime.of(2023, 1, 1, 11, 2, 3)); + assertEquals(ImmutableMap.of("ym", "23-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); + } + + @Test + public void testParsePartitionValues() { + PartitionTimeResolver resolver = + new PartitionTimeResolver( + Arrays.asList("aa", "a", "aaa"), + "$a-$aa-$aaa 00:00:00", + "yyyy-MM-dd HH:mm:ss"); + assertThat(resolver.parsePartitionValues(Arrays.asList("01", "2023", "02"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = + new PartitionTimeResolver( + Arrays.asList("a", "aa", "aaa"), "$aaa$a$aaT000000", "yyyyMMdd'T'HHmmss"); + assertThat(resolver.parsePartitionValues(Arrays.asList("01", "02", "2023"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = + new PartitionTimeResolver(Arrays.asList("aa", "a", "aaa"), "$a$aa$aaa", "yyyyMMdd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("01", "2023", "02"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("dt", "hr"), "$dt $hr", "yyyyMMdd HH"); + assertThat(resolver.parsePartitionValues(Arrays.asList("20230102", "03"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T03:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("hr", "dt"), "$dt $hr", "yyyyMMdd HH"); + assertThat(resolver.parsePartitionValues(Arrays.asList("03", "20230102"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T03:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("dt"), "$dt-$dt", "yyyyMMdd-yyyyMMdd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("20230102"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = + new PartitionTimeResolver( + Arrays.asList("aa", "a", "aaa"), "$a-$aa-$aaa", "yyyy-MM-dd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("01", "2023", "02"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("ymd"), "$ymd", "yyyy-MM-dd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-01-02"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("ymd"), "$ymd", "yyyy-M-d"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-1-2"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("ym"), "$ym-2", "yyyy-M-d"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-1"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("ym"), "$ym-12", "yyyy-M-d"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-1"))) + .isEqualTo(LocalDateTime.parse("2023-01-12T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("y"), "$y-1-2", "yyyy-M-d"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("ym"), "$ym", "yyyy-MM"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023-01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("y"), "$y-12", "yyyy-MM"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023"))) + .isEqualTo(LocalDateTime.parse("2023-12-01T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("y"), "$y-Dec", "yyyy-MMM"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023"))) + .isEqualTo(LocalDateTime.parse("2023-12-01T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("y"), "$y", "yyyy"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + resolver = new PartitionTimeResolver(Arrays.asList("y"), "$y", "yy"); + assertThat(resolver.parsePartitionValues(Arrays.asList("23"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + assertThatThrownBy( + () -> + new PartitionTimeResolver(Arrays.asList("y"), "$y-22", "yyyy-MM") + .parsePartitionValues(Arrays.asList("2026"))) + .satisfies( + PaimonAssertions.anyCauseMatches( + IllegalArgumentException.class, + "Failed to match pattern '$y-22' to formatter 'yyyy-MM'")); + + assertThatThrownBy( + () -> + new PartitionTimeResolver( + Arrays.asList("y", "m"), "$y-$m", "yyyy-MM") + .parsePartitionValues(Arrays.asList("2026", "22"))) + .hasMessageContaining("Text '2026-22' could not be parsed"); + + assertThatThrownBy( + () -> + new PartitionTimeResolver(Arrays.asList("ym"), "$ym", "yyyy-M") + .parsePartitionValues(null)) + .satisfies( + PaimonAssertions.anyCauseMatches( + IllegalArgumentException.class, "Values cannot be null")); + + assertThatThrownBy( + () -> + new PartitionTimeResolver( + Arrays.asList("y", "m"), "$y-$m", "yyyy-M") + .parsePartitionValues(Arrays.asList("2026", "1", "2"))) + .satisfies( + PaimonAssertions.anyCauseMatches( + IllegalArgumentException.class, "Values size mismatch")); + + // Partition columns that are not referenced by the pattern should be ignored. + resolver = new PartitionTimeResolver(Arrays.asList("other", "dt"), "$dt", "yyyy-MM-dd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("dummy", "2023-01-01"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T00:00:00")); + + resolver = + new PartitionTimeResolver( + Arrays.asList("region", "dt", "hour"), + "$dt $hour:00:00", + "yyyy-MM-dd HH:mm:ss"); + assertThat(resolver.parsePartitionValues(Arrays.asList("us", "2023-01-01", "10"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T10:00:00")); + + resolver = + new PartitionTimeResolver( + Arrays.asList("year", "skip", "month", "skip2", "day"), + "$year-$month-$day", + "yyyy-MM-dd"); + assertThat(resolver.parsePartitionValues(Arrays.asList("2023", "x", "01", "y", "02"))) + .isEqualTo(LocalDateTime.parse("2023-01-02T00:00:00")); + + resolver = + new PartitionTimeResolver( + Arrays.asList("hour", "dt", "minute"), + "$dt $hour:$minute:00", + "yyyy-MM-dd HH:mm:ss"); + assertThat(resolver.parsePartitionValues(Arrays.asList("10", "2023-01-01", "30"))) + .isEqualTo(LocalDateTime.parse("2023-01-01T10:30:00")); + } + + @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 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 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"); + } +} 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 50e74d4d4e7f..c29c929ec76e 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.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,6 @@ import java.util.stream.Stream; import static org.assertj.core.api.Assertions.assertThat; -import static org.junit.jupiter.api.Assertions.assertEquals; /** Test class for {@link org.apache.paimon.utils.ChainTableUtils}. */ public class ChainTableUtilsTest { @@ -186,38 +183,6 @@ public void testCreateLinearPredicate() { Assertions.assertTrue(predicate.equals(expected)); } - @Test - public void testGeneratePartitionValues() { - LinkedHashMap 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() { - { - 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() { - { - put("dt", "20230101"); - } - }, - partitionValues); - } - // ========================== Tests for findFirstLatestPartitionsWithProjector // ========================== @@ -917,4 +882,84 @@ private static DataSplit dataSplit(BinaryRow partition, int bucket, List 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 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 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"); + } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneListener.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneListener.java index e83728e38f9a..28b2252398ef 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneListener.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneListener.java @@ -84,7 +84,8 @@ public static Optional create( coreOptions.legacyPartitionName()); PartitionMarkDoneTrigger trigger = - PartitionMarkDoneTrigger.create(coreOptions, isRestored, stateStore); + PartitionMarkDoneTrigger.create( + coreOptions, table.partitionKeys(), isRestored, stateStore); List actions = PartitionMarkDoneAction.createActions(cl, table, coreOptions); diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTrigger.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTrigger.java index 8ddbadb0bc9f..55bdeb9fee30 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTrigger.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTrigger.java @@ -23,7 +23,7 @@ import org.apache.paimon.flink.sink.state.StateStore; import org.apache.paimon.fs.Path; import org.apache.paimon.options.Options; -import org.apache.paimon.partition.PartitionTimeExtractor; +import org.apache.paimon.partition.PartitionTimeResolvable; import org.apache.paimon.utils.StringUtils; import org.apache.flink.api.common.state.ListState; @@ -43,6 +43,7 @@ import java.util.Collections; import java.util.HashMap; import java.util.Iterator; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.Optional; @@ -62,7 +63,7 @@ public class PartitionMarkDoneTrigger { new ListSerializer<>(StringSerializer.INSTANCE)); private final State state; - private final PartitionTimeExtractor timeExtractor; + private final PartitionTimeResolvable timeResolver; // can be null when markDoneWhenEndInput is true @Nullable private final Long timeInterval; // can be null when markDoneWhenEndInput is true @@ -72,14 +73,14 @@ public class PartitionMarkDoneTrigger { public PartitionMarkDoneTrigger( State state, - PartitionTimeExtractor timeExtractor, + PartitionTimeResolvable timeResolver, @Nullable Duration timeInterval, @Nullable Duration idleTime, boolean markDoneWhenEndInput) throws Exception { this( state, - timeExtractor, + timeResolver, timeInterval, idleTime, System.currentTimeMillis(), @@ -88,7 +89,7 @@ public PartitionMarkDoneTrigger( public PartitionMarkDoneTrigger( State state, - PartitionTimeExtractor timeExtractor, + PartitionTimeResolvable timeResolver, @Nullable Duration timeInterval, @Nullable Duration idleTime, long currentTimeMillis, @@ -96,7 +97,7 @@ public PartitionMarkDoneTrigger( throws Exception { this.pendingPartitions = new HashMap<>(); this.state = state; - this.timeExtractor = timeExtractor; + this.timeResolver = timeResolver; this.timeInterval = timeInterval == null ? null : timeInterval.toMillis(); this.idleTime = idleTime == null ? null : idleTime.toMillis(); this.markDoneWhenEndInput = markDoneWhenEndInput; @@ -204,8 +205,13 @@ List donePartitions( @VisibleForTesting Optional extractDateTime(String partition) { try { - return Optional.of( - timeExtractor.extract(extractPartitionSpecFromPath(new Path(partition)))); + LinkedHashMap spec = extractPartitionSpecFromPath(new Path(partition)); + List partitionKeys = timeResolver.partitionKeys(); + List values = new ArrayList<>(partitionKeys.size()); + for (String key : partitionKeys) { + values.add(spec.get(key)); + } + return Optional.of(timeResolver.parsePartitionValues(values)); } catch (DateTimeParseException e) { LOG.warn( "Can't extract datetime from partition {}, please check configuration items 'partition.timestamp-formatter' and 'partition.timestamp-pattern'.", @@ -256,11 +262,16 @@ public void update(List partitions) throws Exception { } public static PartitionMarkDoneTrigger create( - CoreOptions coreOptions, boolean isRestored, StateStore stateStore) throws Exception { + CoreOptions coreOptions, + List partitionKeys, + boolean isRestored, + StateStore stateStore) + throws Exception { Options options = coreOptions.toConfiguration(); return new PartitionMarkDoneTrigger( new PartitionMarkDoneTrigger.PartitionMarkDoneTriggerState(isRestored, stateStore), - new PartitionTimeExtractor( + PartitionTimeResolvable.create( + partitionKeys, coreOptions.partitionTimestampPattern(), coreOptions.partitionTimestampFormatter()), options.get(PARTITION_TIME_INTERVAL), diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTriggerTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTriggerTest.java index 84914253fc14..cc52bbe2c1a8 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTriggerTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/listener/PartitionMarkDoneTriggerTest.java @@ -18,7 +18,7 @@ package org.apache.paimon.flink.sink.listener; -import org.apache.paimon.partition.PartitionTimeExtractor; +import org.apache.paimon.partition.PartitionTimeResolver; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; @@ -29,6 +29,7 @@ import java.time.LocalTime; import java.time.ZoneId; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import static org.assertj.core.api.Assertions.assertThat; @@ -40,7 +41,7 @@ class PartitionMarkDoneTriggerTest { private List pendingPartitions; private PartitionMarkDoneTrigger.State state; - private PartitionTimeExtractor extractor; + private PartitionTimeResolver resolver; @BeforeEach public void before() throws Exception { @@ -58,7 +59,8 @@ public void update(List partitions) { pendingPartitions.addAll(partitions); } }; - this.extractor = new PartitionTimeExtractor("$dt", "yyyy-MM-dd"); + this.resolver = + new PartitionTimeResolver(Collections.singletonList("dt"), "$dt", "yyyy-MM-dd"); } @Test @@ -66,7 +68,7 @@ public void testWithoutEndInput() throws Exception { PartitionMarkDoneTrigger trigger = new PartitionMarkDoneTrigger( state, - extractor, + resolver, timeInterval, idleTime, toEpochMillis("2024-02-01"), @@ -112,7 +114,7 @@ public void testWithoutEndInput() throws Exception { trigger = new PartitionMarkDoneTrigger( state, - extractor, + resolver, timeInterval, idleTime, toEpochMillis("2024-02-06"), @@ -129,12 +131,7 @@ public void testWithoutEndInput() throws Exception { public void testWithEndInput() throws Exception { PartitionMarkDoneTrigger trigger = new PartitionMarkDoneTrigger( - state, - extractor, - timeInterval, - idleTime, - toEpochMillis("2024-02-01"), - true); + state, resolver, timeInterval, idleTime, toEpochMillis("2024-02-01"), true); // test not reach partition end + idle time trigger.notifyPartition("dt=2024-02-02", toEpochMillis("2024-02-01")); @@ -146,12 +143,7 @@ public void testWithEndInput() throws Exception { public void testParseNonDateFormattedPartition() throws Exception { PartitionMarkDoneTrigger trigger = new PartitionMarkDoneTrigger( - state, - extractor, - timeInterval, - idleTime, - toEpochMillis("2024-02-01"), - true); + state, resolver, timeInterval, idleTime, toEpochMillis("2024-02-01"), true); assertThat(trigger.extractDateTime("unknown")).isEmpty(); trigger.notifyPartition("dt=__DEFAULT_PARTITION__", toEpochMillis("2024-02-01")); 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 ed3814ab7824..1b6bff959ede 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 void testChainTableWithMultiChainKeys(@TempDir java.nio.file.Path tempDir 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 {