Skip to content

Commit 235a52a

Browse files
author
张文领
committed
Add Maintenance tab for table process management
1 parent be6663b commit 235a52a

17 files changed

Lines changed: 234 additions & 33 deletions

File tree

amoro-ams/src/main/java/org/apache/amoro/server/dashboard/DashboardServer.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -270,6 +270,9 @@ private EndpointGroup apiGroup() {
270270
get(
271271
"/catalogs/{catalog}/dbs/{db}/tables/{table}/optimizing-types",
272272
tableController::getOptimizingTypes);
273+
get(
274+
"/catalogs/{catalog}/dbs/{db}/tables/{table}/maintenance-types",
275+
tableController::getMaintenanceTypes);
273276
get(
274277
"/catalogs/{catalog}/dbs/{db}/tables/{table}/optimizing-processes/{processId}/tasks",
275278
tableController::getOptimizingProcessTasks);

amoro-ams/src/main/java/org/apache/amoro/server/dashboard/MixedAndIcebergTableDescriptor.java

Lines changed: 50 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import com.github.pagehelper.PageHelper;
2323
import com.github.pagehelper.PageInfo;
2424
import org.apache.amoro.AmoroTable;
25+
import org.apache.amoro.IcebergActions;
2526
import org.apache.amoro.ServerTableIdentifier;
2627
import org.apache.amoro.TableFormat;
2728
import org.apache.amoro.api.CommitMetaProducer;
@@ -73,6 +74,7 @@
7374
import org.apache.amoro.utils.MixedDataFiles;
7475
import org.apache.amoro.utils.MixedTableUtil;
7576
import org.apache.commons.collections.CollectionUtils;
77+
import org.apache.commons.lang3.StringUtils;
7678
import org.apache.commons.lang3.tuple.Pair;
7779
import org.apache.iceberg.ContentFile;
7880
import org.apache.iceberg.FileScanTask;
@@ -118,6 +120,12 @@ public class MixedAndIcebergTableDescriptor extends PersistentBase
118120

119121
private static final Logger LOG = LoggerFactory.getLogger(MixedAndIcebergTableDescriptor.class);
120122

123+
private static final List<String> OPTIMIZING_TYPE_LIST =
124+
Arrays.stream(OptimizingType.values()).map(Enum::name).collect(Collectors.toList());
125+
126+
private static final String OPTIMIZING = "OPTIMIZING";
127+
private static final String MAINTENANCE = "MAINTENANCE";
128+
121129
private ExecutorService executorService;
122130

123131
@Override
@@ -655,7 +663,12 @@ public List<ConsumerInfo> getTableConsumerInfos(AmoroTable<?> amoroTable) {
655663

656664
@Override
657665
public Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
658-
AmoroTable<?> amoroTable, String type, ProcessStatus status, int limit, int offset) {
666+
AmoroTable<?> amoroTable,
667+
String type,
668+
String processCategory,
669+
ProcessStatus status,
670+
int limit,
671+
int offset) {
659672
TableIdentifier tableIdentifier = amoroTable.id();
660673
ServerTableIdentifier identifier =
661674
getAs(
@@ -671,12 +684,35 @@ public Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
671684
int total = 0;
672685
// page helper is 1-based
673686
int pageNumber = (offset / limit) + 1;
687+
688+
// Only apply category filtering when type is not specified
689+
final List<String> includeTypes;
690+
final List<String> excludeTypes;
691+
692+
if (StringUtils.isBlank(type)) {
693+
if (OPTIMIZING.equalsIgnoreCase(processCategory)) {
694+
includeTypes = OPTIMIZING_TYPE_LIST;
695+
excludeTypes = null;
696+
} else if (MAINTENANCE.equalsIgnoreCase(processCategory)) {
697+
includeTypes = null;
698+
excludeTypes = OPTIMIZING_TYPE_LIST;
699+
} else {
700+
includeTypes = null;
701+
excludeTypes = null;
702+
}
703+
} else {
704+
includeTypes = null;
705+
excludeTypes = null;
706+
}
707+
674708
List<TableProcessMeta> processMetaList = Collections.emptyList();
675709
try (Page<?> ignored = PageHelper.startPage(pageNumber, limit, true)) {
676710
processMetaList =
677711
getAs(
678712
TableProcessMapper.class,
679-
mapper -> mapper.listProcessMeta(identifier.getId(), type, status));
713+
mapper ->
714+
mapper.listProcessMeta(
715+
identifier.getId(), type, includeTypes, excludeTypes, status));
680716
PageInfo<TableProcessMeta> pageInfo = new PageInfo<>(processMetaList);
681717
total = (int) pageInfo.getTotal();
682718
LOG.info(
@@ -716,6 +752,18 @@ public Map<String, String> getTableOptimizingTypes(AmoroTable<?> amoroTable) {
716752
return types;
717753
}
718754

755+
@Override
756+
public Map<String, String> getTableMaintenanceTypes(AmoroTable<?> amoroTable) {
757+
Map<String, String> types = Maps.newHashMap();
758+
types.put(IcebergActions.EXPIRE_SNAPSHOTS.getName(), "Expire Snapshots");
759+
types.put(IcebergActions.CLEAN_ORPHAN.getName(), "Clean Orphan Files");
760+
types.put(IcebergActions.CLEAN_DANGLING_DELETE.getName(), "Clean Dangling Delete Files");
761+
types.put(IcebergActions.EXPIRE_DATA.getName(), "Expire Data");
762+
types.put(IcebergActions.SYNC_HIVE_TABLES.getName(), "Sync Hive Tables");
763+
types.put(IcebergActions.AUTO_CREATE_TAGS.getName(), "Auto Create Tags");
764+
return types;
765+
}
766+
719767
@Override
720768
public List<OptimizingTaskInfo> getOptimizingTaskInfos(
721769
AmoroTable<?> amoroTable, String processId) {

amoro-ams/src/main/java/org/apache/amoro/server/dashboard/ServerTableDescriptor.java

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -129,11 +129,16 @@ public List<ConsumerInfo> getTableConsumersInfos(TableIdentifier tableIdentifier
129129
}
130130

131131
public Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
132-
TableIdentifier tableIdentifier, String type, ProcessStatus status, int limit, int offset) {
132+
TableIdentifier tableIdentifier,
133+
String type,
134+
String processCategory,
135+
ProcessStatus status,
136+
int limit,
137+
int offset) {
133138
AmoroTable<?> amoroTable = loadTable(tableIdentifier);
134139
FormatTableDescriptor formatTableDescriptor = formatDescriptorMap.get(amoroTable.format());
135140
return formatTableDescriptor.getOptimizingProcessesInfo(
136-
amoroTable, type, status, limit, offset);
141+
amoroTable, type, processCategory, status, limit, offset);
137142
}
138143

139144
public List<OptimizingTaskInfo> getOptimizingProcessTaskInfos(
@@ -149,6 +154,12 @@ public Map<String, String> getTableOptimizingTypes(TableIdentifier tableIdentifi
149154
return formatTableDescriptor.getTableOptimizingTypes(amoroTable);
150155
}
151156

157+
public Map<String, String> getTableMaintenanceTypes(TableIdentifier tableIdentifier) {
158+
AmoroTable<?> amoroTable = loadTable(tableIdentifier);
159+
FormatTableDescriptor formatTableDescriptor = formatDescriptorMap.get(amoroTable.format());
160+
return formatTableDescriptor.getTableMaintenanceTypes(amoroTable);
161+
}
162+
152163
private AmoroTable<?> loadTable(TableIdentifier identifier) {
153164
ServerCatalog catalog = catalogManager.getServerCatalog(identifier.getCatalog());
154165
return catalog.loadTable(identifier.getDatabase(), identifier.getTableName());

amoro-ams/src/main/java/org/apache/amoro/server/dashboard/controller/TableController.java

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -332,6 +332,11 @@ public void getOptimizingProcesses(Context ctx) {
332332
type = null;
333333
}
334334

335+
String processCategory = ctx.queryParam("processCategory");
336+
if (StringUtils.isBlank(processCategory)) {
337+
processCategory = null;
338+
}
339+
335340
String status = ctx.queryParam("status");
336341
Integer page = ctx.queryParamAsClass("page", Integer.class).getOrDefault(1);
337342
Integer pageSize = ctx.queryParamAsClass("pageSize", Integer.class).getOrDefault(20);
@@ -346,7 +351,12 @@ public void getOptimizingProcesses(Context ctx) {
346351
StringUtils.isBlank(status) ? null : ProcessStatus.valueOf(status);
347352
Pair<List<OptimizingProcessInfo>, Integer> optimizingProcessesInfo =
348353
tableDescriptor.getOptimizingProcessesInfo(
349-
tableIdentifier.buildTableIdentifier(), type, processStatus, limit, offset);
354+
tableIdentifier.buildTableIdentifier(),
355+
type,
356+
processCategory,
357+
processStatus,
358+
limit,
359+
offset);
350360
List<OptimizingProcessInfo> result = optimizingProcessesInfo.getLeft();
351361
int total = optimizingProcessesInfo.getRight();
352362

@@ -364,6 +374,17 @@ public void getOptimizingTypes(Context ctx) {
364374
ctx.json(OkResponse.of(values));
365375
}
366376

377+
public void getMaintenanceTypes(Context ctx) {
378+
String catalog = ctx.pathParam("catalog");
379+
String db = ctx.pathParam("db");
380+
String table = ctx.pathParam("table");
381+
TableIdentifier tableIdentifier = TableIdentifier.of(catalog, db, table);
382+
383+
Map<String, String> values =
384+
tableDescriptor.getTableMaintenanceTypes(tableIdentifier.buildTableIdentifier());
385+
ctx.json(OkResponse.of(values));
386+
}
387+
367388
/**
368389
* Get tasks of optimizing process.
369390
*

amoro-ams/src/main/java/org/apache/amoro/server/persistence/mapper/TableProcessMapper.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -131,13 +131,17 @@ void updateProcess(
131131
+ "create_time, finish_time, fail_message, process_parameters, summary "
132132
+ "FROM table_process WHERE table_id = #{tableId} "
133133
+ " <if test='processType != null'> AND process_type = #{processType}</if>"
134+
+ " <if test='processType == null and includeTypes != null and includeTypes.size() > 0'> AND process_type IN <foreach collection='includeTypes' item='type' open='(' separator=',' close=')'>#{type}</foreach></if>"
135+
+ " <if test='processType == null and excludeTypes != null and excludeTypes.size() > 0'> AND process_type NOT IN <foreach collection='excludeTypes' item='type' open='(' separator=',' close=')'>#{type}</foreach></if>"
134136
+ " <if test='status != null'> AND status = #{status}</if>"
135137
+ " ORDER BY process_id desc"
136138
+ "</script>")
137139
@ResultMap("tableProcessMap")
138140
List<TableProcessMeta> listProcessMeta(
139141
@Param("tableId") long tableId,
140142
@Param("processType") String processType,
143+
@Param("includeTypes") List<String> includeTypes,
144+
@Param("excludeTypes") List<String> excludeTypes,
141145
@Param("status") ProcessStatus optimizingStatus);
142146

143147
@Select(

amoro-ams/src/test/java/org/apache/amoro/server/dashboard/TestIcebergServerTableDescriptor.java

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -301,21 +301,22 @@ public void testOptimizingProcess() {
301301
doReturn(tableIdentifier).when(table).id();
302302

303303
Pair<List<OptimizingProcessInfo>, Integer> res =
304-
descriptor.getOptimizingProcessesInfo(table, null, null, 4, 4);
304+
descriptor.getOptimizingProcessesInfo(table, null, null, null, 4, 4);
305305
Integer expectReturnItemSizeForNoTypeNoStatusOffset0Limit5 = 4;
306306
Integer expectTotalForNoTypeNoStatusOffset0Limit5 = 10;
307307
Assert.assertEquals(
308308
expectReturnItemSizeForNoTypeNoStatusOffset0Limit5, (Integer) res.getLeft().size());
309309
Assert.assertEquals(expectTotalForNoTypeNoStatusOffset0Limit5, res.getRight());
310310

311-
res = descriptor.getOptimizingProcessesInfo(table, null, ProcessStatus.SUCCESS, 5, 0);
311+
res = descriptor.getOptimizingProcessesInfo(table, null, null, ProcessStatus.SUCCESS, 5, 0);
312312
Integer expectReturnItemSizeForOnlyStatusOffset0limit5 = 5;
313313
Integer expectedTotalForOnlyStatusOffset0Limit5 = 7;
314314
Assert.assertEquals(
315315
expectReturnItemSizeForOnlyStatusOffset0limit5, (Integer) res.getLeft().size());
316316
Assert.assertEquals(expectedTotalForOnlyStatusOffset0Limit5, res.getRight());
317317

318-
res = descriptor.getOptimizingProcessesInfo(table, OptimizingType.MINOR.name(), null, 5, 0);
318+
res =
319+
descriptor.getOptimizingProcessesInfo(table, OptimizingType.MINOR.name(), null, null, 5, 0);
319320
Integer expectedRetItemsSizeForOnlyTypeOffset0Limit5 = 4;
320321
Integer expectedRetTotalForOnlyTypeOffset0Limit5 = 4;
321322
Assert.assertEquals(
@@ -324,7 +325,7 @@ public void testOptimizingProcess() {
324325

325326
res =
326327
descriptor.getOptimizingProcessesInfo(
327-
table, OptimizingType.MINOR.name(), ProcessStatus.SUCCESS, 2, 2);
328+
table, OptimizingType.MINOR.name(), null, ProcessStatus.SUCCESS, 2, 2);
328329
Integer expectedRetItemSizeForBothTypeAndStatusOffset2Limit2 = 2;
329330
Integer expectedRetTotalForBothTypeAndStatusOffset2Limit2 = 4;
330331
Assert.assertEquals(

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

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,8 @@ protected TableProcessMeta waitOptimizeResult() {
124124
List<TableProcessMeta> tableOptimizingProcesses =
125125
getAs(
126126
TableProcessMapper.class,
127-
mapper -> mapper.listProcessMeta(identifier.getId(), null, null));
127+
mapper ->
128+
mapper.listProcessMeta(identifier.getId(), null, null, null, null));
128129
if (tableOptimizingProcesses == null || tableOptimizingProcesses.isEmpty()) {
129130
LOG.info("optimize history is empty");
130131
return Status.RUNNING;
@@ -157,7 +158,7 @@ protected TableProcessMeta waitOptimizeResult() {
157158
List<TableProcessMeta> result =
158159
getAs(
159160
TableProcessMapper.class,
160-
mapper -> mapper.listProcessMeta(identifier.getId(), null, null))
161+
mapper -> mapper.listProcessMeta(identifier.getId(), null, null, null, null))
161162
.stream()
162163
.filter(p -> p.getProcessId() > lastProcessId)
163164
.filter(p -> p.getStatus().equals(ProcessStatus.SUCCESS))
@@ -190,7 +191,7 @@ protected void assertOptimizeHangUp() {
190191
List<TableProcessMeta> tableOptimizingProcesses =
191192
getAs(
192193
TableProcessMapper.class,
193-
mapper -> mapper.listProcessMeta(identifier.getId(), null, null))
194+
mapper -> mapper.listProcessMeta(identifier.getId(), null, null, null, null))
194195
.stream()
195196
.filter(p -> p.getProcessId() > lastProcessId)
196197
.collect(Collectors.toList());

amoro-ams/src/test/java/org/apache/amoro/server/scheduler/inline/TestProcessDataExpiringExecutor.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -184,7 +184,9 @@ public void insertProcess(long tableId, long processId, ProcessStatus status, lo
184184
}
185185

186186
public List<TableProcessMeta> listProcesses(long tableId) {
187-
return getAs(TableProcessMapper.class, mapper -> mapper.listProcessMeta(tableId, null, null));
187+
return getAs(
188+
TableProcessMapper.class,
189+
mapper -> mapper.listProcessMeta(tableId, null, null, null, null));
188190
}
189191

190192
public void cleanAll(long tableId) {

amoro-common/src/main/java/org/apache/amoro/table/descriptor/FormatTableDescriptor.java

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,8 +23,10 @@
2323
import org.apache.amoro.process.ProcessStatus;
2424
import org.apache.commons.lang3.tuple.Pair;
2525

26+
import java.util.Collections;
2627
import java.util.List;
2728
import java.util.Map;
29+
import java.util.Set;
2830
import java.util.concurrent.ExecutorService;
2931

3032
/** API for obtaining metadata information of various formats. */
@@ -63,11 +65,37 @@ List<PartitionFileBaseInfo> getTableFiles(
6365

6466
/** Get the paged optimizing process information of the {@link AmoroTable} and total size. */
6567
Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
66-
AmoroTable<?> amoroTable, String type, ProcessStatus status, int limit, int offset);
68+
AmoroTable<?> amoroTable,
69+
String type,
70+
String processCategory,
71+
ProcessStatus status,
72+
int limit,
73+
int offset);
6774

6875
/** Return the optimizing types of the {@link AmoroTable} is supported. */
6976
Map<String, String> getTableOptimizingTypes(AmoroTable<?> amoroTable);
7077

78+
/** Return the maintenance types of the {@link AmoroTable} is supported. */
79+
default Map<String, String> getTableMaintenanceTypes(AmoroTable<?> amoroTable) {
80+
return Collections.emptyMap();
81+
}
82+
83+
static boolean matchProcessCategory(
84+
String processCategory, Set<String> optimizingTypes, String type) {
85+
if (processCategory == null) {
86+
return true;
87+
}
88+
boolean isOptimizingType =
89+
type != null && optimizingTypes.stream().anyMatch(t -> t.equalsIgnoreCase(type));
90+
if ("OPTIMIZING".equalsIgnoreCase(processCategory)) {
91+
return isOptimizingType;
92+
}
93+
if ("MAINTENANCE".equalsIgnoreCase(processCategory)) {
94+
return !isOptimizingType;
95+
}
96+
return true;
97+
}
98+
7199
/** Get the paged optimizing process tasks information of the {@link AmoroTable}. */
72100
List<OptimizingTaskInfo> getOptimizingTaskInfos(AmoroTable<?> amoroTable, String processId);
73101

amoro-format-hudi/src/main/java/org/apache/amoro/formats/hudi/HudiTableDescriptor.java

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,7 @@ public class HudiTableDescriptor implements FormatTableDescriptor {
9797
private static final Logger LOG = LoggerFactory.getLogger(HudiTableDescriptor.class);
9898
private static final String COMPACTION = "compaction";
9999
private static final String CLUSTERING = "clustering";
100+
private static final Set<String> OPTIMIZING_TYPES = Sets.newHashSet(COMPACTION, CLUSTERING);
100101
// table comment
101102
private static final String COMMENT = "comment";
102103

@@ -340,7 +341,12 @@ private Stream<PartitionFileBaseInfo> fileSliceToFileStream(String partition, Fi
340341

341342
@Override
342343
public Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
343-
AmoroTable<?> amoroTable, String type, ProcessStatus status, int limit, int offset) {
344+
AmoroTable<?> amoroTable,
345+
String type,
346+
String processCategory,
347+
ProcessStatus status,
348+
int limit,
349+
int offset) {
344350
HoodieJavaTable hoodieTable = (HoodieJavaTable) amoroTable.originalTable();
345351
HoodieDefaultTimeline timeline = new HoodieActiveTimeline(hoodieTable.getMetaClient(), false);
346352
List<HoodieInstant> instants = timeline.getInstants();
@@ -387,6 +393,10 @@ public Pair<List<OptimizingProcessInfo>, Integer> getOptimizingProcessesInfo(
387393
i ->
388394
StringUtils.isNullOrEmpty(type) || type.equalsIgnoreCase(i.getOptimizingType()))
389395
.filter(i -> status == null || status == i.getStatus())
396+
.filter(
397+
i ->
398+
FormatTableDescriptor.matchProcessCategory(
399+
processCategory, OPTIMIZING_TYPES, i.getOptimizingType()))
390400
.collect(Collectors.toList());
391401
int total = infos.size();
392402
infos = infos.stream().skip(offset).limit(limit).collect(Collectors.toList());

0 commit comments

Comments
 (0)