Skip to content

Commit 2ff66e5

Browse files
lintingbinzhoujinsongklion26xxubaibaiyangtx
authored
[AMORO-3272] data-expire by partition info (#3273)
* feature: data-expire by partition info * add comments * Refactor the code to improve readability. * Refactor the code to improve readability. * feature: add test cases * fix test error * optimize code * update * checkstyle --------- Co-authored-by: ZhouJinsong <zhoujinsong0505@163.com> Co-authored-by: Congxian Qiu <qcx978132955@gmail.com> Co-authored-by: Xavier Bai <xuba@apache.org> Co-authored-by: baiyangtx <xiangnebula@163.com> Co-authored-by: Xavier Bai <xuba@cisco.com>
1 parent b81d3e1 commit 2ff66e5

3 files changed

Lines changed: 97 additions & 8 deletions

File tree

amoro-ams/src/main/java/org/apache/amoro/server/optimizing/maintainer/IcebergTableMaintainer.java

Lines changed: 78 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,8 @@
4545
import org.apache.iceberg.DeleteFiles;
4646
import org.apache.iceberg.FileContent;
4747
import org.apache.iceberg.FileScanTask;
48+
import org.apache.iceberg.PartitionField;
49+
import org.apache.iceberg.PartitionSpec;
4850
import org.apache.iceberg.ReachableFileUtil;
4951
import org.apache.iceberg.RewriteFiles;
5052
import org.apache.iceberg.Schema;
@@ -63,6 +65,7 @@
6365
import org.apache.iceberg.types.Conversions;
6466
import org.apache.iceberg.types.Type;
6567
import org.apache.iceberg.types.Types;
68+
import org.apache.iceberg.util.SerializableFunction;
6669
import org.apache.iceberg.util.ThreadPools;
6770
import org.slf4j.Logger;
6871
import org.slf4j.LoggerFactory;
@@ -75,8 +78,10 @@
7578
import java.time.ZoneId;
7679
import java.time.ZoneOffset;
7780
import java.time.format.DateTimeFormatter;
81+
import java.util.ArrayList;
7882
import java.util.Collections;
7983
import java.util.HashSet;
84+
import java.util.List;
8085
import java.util.Locale;
8186
import java.util.Map;
8287
import java.util.Optional;
@@ -650,7 +655,10 @@ private Set<String> deleteInvalidMetadataFile(
650655
}
651656

652657
CloseableIterable<FileEntry> fileScan(
653-
Table table, Expression dataFilter, DataExpirationConfig expirationConfig) {
658+
Table table,
659+
Expression dataFilter,
660+
DataExpirationConfig expirationConfig,
661+
long expireTimestamp) {
654662
TableScan tableScan = table.newScan().filter(dataFilter).includeColumnStats();
655663

656664
CloseableIterable<FileScanTask> tasks;
@@ -680,6 +688,7 @@ CloseableIterable<FileEntry> fileScan(
680688
.collect(Collectors.toSet());
681689

682690
Types.NestedField field = table.schema().findField(expirationConfig.getExpirationField());
691+
Comparable<?> expireValue = getExpireValue(expirationConfig, field, expireTimestamp);
683692
return CloseableIterable.transform(
684693
CloseableIterable.withNoopClose(Iterables.concat(dataFiles, deleteFiles)),
685694
contentFile -> {
@@ -689,7 +698,8 @@ CloseableIterable<FileEntry> fileScan(
689698
field,
690699
DateTimeFormatter.ofPattern(
691700
expirationConfig.getDateTimePattern(), Locale.getDefault()),
692-
expirationConfig.getNumberDateFormat());
701+
expirationConfig.getNumberDateFormat(),
702+
expireValue);
693703
return new FileEntry(contentFile.copyWithoutStats(), literal);
694704
});
695705
}
@@ -698,7 +708,8 @@ protected ExpireFiles expiredFileScan(
698708
DataExpirationConfig expirationConfig, Expression dataFilter, long expireTimestamp) {
699709
Map<StructLike, DataFileFreshness> partitionFreshness = Maps.newConcurrentMap();
700710
ExpireFiles expiredFiles = new ExpireFiles();
701-
try (CloseableIterable<FileEntry> entries = fileScan(table, dataFilter, expirationConfig)) {
711+
try (CloseableIterable<FileEntry> entries =
712+
fileScan(table, dataFilter, expirationConfig, expireTimestamp)) {
702713
Queue<FileEntry> fileEntries = new LinkedTransferQueue<>();
703714
entries.forEach(
704715
e -> {
@@ -716,6 +727,33 @@ protected ExpireFiles expiredFileScan(
716727
return expiredFiles;
717728
}
718729

730+
private Comparable<?> getExpireValue(
731+
DataExpirationConfig expirationConfig, Types.NestedField field, long expireTimestamp) {
732+
switch (field.type().typeId()) {
733+
// expireTimestamp is in milliseconds, TIMESTAMP type is in microseconds
734+
case TIMESTAMP:
735+
return expireTimestamp * 1000;
736+
case LONG:
737+
if (expirationConfig.getNumberDateFormat().equals(EXPIRE_TIMESTAMP_MS)) {
738+
return expireTimestamp;
739+
} else if (expirationConfig.getNumberDateFormat().equals(EXPIRE_TIMESTAMP_S)) {
740+
return expireTimestamp / 1000;
741+
} else {
742+
throw new IllegalArgumentException(
743+
"Number dateformat: " + expirationConfig.getNumberDateFormat());
744+
}
745+
case STRING:
746+
return LocalDateTime.ofInstant(
747+
Instant.ofEpochMilli(expireTimestamp), getDefaultZoneId(field))
748+
.format(
749+
DateTimeFormatter.ofPattern(
750+
expirationConfig.getDateTimePattern(), Locale.getDefault()));
751+
default:
752+
throw new IllegalArgumentException(
753+
"Unsupported expiration field type: " + field.type().typeId());
754+
}
755+
}
756+
719757
/**
720758
* Create a filter expression for expired files for the `FILE` level. For the `PARTITION` level,
721759
* we need to collect the oldest files to determine if the partition is obsolete, so we will not
@@ -917,17 +955,20 @@ static boolean willNotRetain(
917955
}
918956
}
919957

920-
private static Literal<Long> getExpireTimestampLiteral(
958+
private Literal<Long> getExpireTimestampLiteral(
921959
ContentFile<?> contentFile,
922960
Types.NestedField field,
923961
DateTimeFormatter formatter,
924-
String numberDateFormatter) {
962+
String numberDateFormatter,
963+
Comparable<?> expireValue) {
925964
Type type = field.type();
926965
Object upperBound =
927966
Conversions.fromByteBuffer(type, contentFile.upperBounds().get(field.fieldId()));
928967
Literal<Long> literal = Literal.of(Long.MAX_VALUE);
929968
if (null == upperBound) {
930-
return literal;
969+
if (canBeExpireByPartitionValue(contentFile, field, expireValue)) {
970+
literal = Literal.of(0L);
971+
}
931972
} else if (upperBound instanceof Long) {
932973
switch (type.typeId()) {
933974
case TIMESTAMP:
@@ -951,9 +992,40 @@ private static Literal<Long> getExpireTimestampLiteral(
951992
.toInstant()
952993
.toEpochMilli());
953994
}
995+
954996
return literal;
955997
}
956998

999+
@SuppressWarnings("unchecked")
1000+
private boolean canBeExpireByPartitionValue(
1001+
ContentFile<?> contentFile, Types.NestedField expireField, Comparable<?> expireValue) {
1002+
PartitionSpec partitionSpec = table.specs().get(contentFile.specId());
1003+
int pos = 0;
1004+
List<Boolean> compareResults = new ArrayList<>();
1005+
for (PartitionField partitionField : partitionSpec.fields()) {
1006+
if (partitionField.sourceId() == expireField.fieldId()) {
1007+
if (partitionField.transform().isVoid()) {
1008+
return false;
1009+
}
1010+
1011+
Comparable<?> partitionUpperBound =
1012+
((SerializableFunction<Comparable<?>, Comparable<?>>)
1013+
partitionField.transform().bind(expireField.type()))
1014+
.apply(expireValue);
1015+
Comparable<Object> filePartitionValue =
1016+
contentFile.partition().get(pos, partitionUpperBound.getClass());
1017+
int compared = filePartitionValue.compareTo(partitionUpperBound);
1018+
Boolean compareResult =
1019+
expireField.type() == Types.StringType.get() ? compared <= 0 : compared < 0;
1020+
compareResults.add(compareResult);
1021+
}
1022+
1023+
pos++;
1024+
}
1025+
1026+
return !compareResults.isEmpty() && compareResults.stream().allMatch(Boolean::booleanValue);
1027+
}
1028+
9571029
public Table getTable() {
9581030
return table;
9591031
}

amoro-ams/src/main/java/org/apache/amoro/server/optimizing/maintainer/MixedTableMaintainer.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -195,11 +195,11 @@ public void expireDataFrom(DataExpirationConfig expirationConfig, Instant instan
195195

196196
CloseableIterable<MixedFileEntry> changeEntries =
197197
CloseableIterable.transform(
198-
changeMaintainer.fileScan(changeTable, dataFilter, expirationConfig),
198+
changeMaintainer.fileScan(changeTable, dataFilter, expirationConfig, expireTimestamp),
199199
e -> new MixedFileEntry(e.getFile(), e.getTsBound(), true));
200200
CloseableIterable<MixedFileEntry> baseEntries =
201201
CloseableIterable.transform(
202-
baseMaintainer.fileScan(baseTable, dataFilter, expirationConfig),
202+
baseMaintainer.fileScan(baseTable, dataFilter, expirationConfig, expireTimestamp),
203203
e -> new MixedFileEntry(e.getFile(), e.getTsBound(), false));
204204
IcebergTableMaintainer.ExpireFiles changeExpiredFiles =
205205
new IcebergTableMaintainer.ExpireFiles();

amoro-ams/src/test/java/org/apache/amoro/server/optimizing/maintainer/TestDataExpire.java

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020

2121
import static org.apache.amoro.BasicTableTestHelper.PRIMARY_KEY_SPEC;
2222
import static org.apache.amoro.BasicTableTestHelper.SPEC;
23+
import static org.junit.Assume.assumeTrue;
2324

2425
import org.apache.amoro.BasicTableTestHelper;
2526
import org.apache.amoro.TableFormat;
@@ -49,6 +50,7 @@
4950
import org.apache.commons.lang.StringUtils;
5051
import org.apache.iceberg.ContentFile;
5152
import org.apache.iceberg.DeleteFile;
53+
import org.apache.iceberg.MetricsModes;
5254
import org.apache.iceberg.PartitionSpec;
5355
import org.apache.iceberg.Schema;
5456
import org.apache.iceberg.Table;
@@ -194,6 +196,7 @@ private void testUnKeyedPartitionLevel() {
194196

195197
List<Record> expected;
196198
if (tableTestHelper().partitionSpec().isPartitioned()) {
199+
// retention time is 1 day, expire partitions that order than 2022-01-02
197200
if (expireByStringDate()) {
198201
expected =
199202
Lists.newArrayList(
@@ -464,6 +467,20 @@ public void testNormalFieldFileLevel() {
464467
testFileLevel();
465468
}
466469

470+
@Test
471+
public void testExpireByPartitionWhenMetricsModeIsNone() {
472+
assumeTrue(getMixedTable().format().in(TableFormat.MIXED_ICEBERG, TableFormat.ICEBERG));
473+
474+
getMixedTable()
475+
.updateProperties()
476+
.set(
477+
org.apache.iceberg.TableProperties.DEFAULT_WRITE_METRICS_MODE,
478+
MetricsModes.None.get().toString())
479+
.commit();
480+
481+
testPartitionLevel();
482+
}
483+
467484
@Test
468485
public void testGcDisabled() {
469486
MixedTable testTable = getMixedTable();

0 commit comments

Comments
 (0)