BanyanDBMetricsDAO: handle storeIDTag in multiGet for BanyanDBModelExtension (#11002)

This commit is contained in:
Gao Hongtao 2023-06-27 15:40:20 +08:00 committed by GitHub
parent c5c268b653
commit 5ad74eb261
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
2 changed files with 38 additions and 19 deletions

View File

@ -25,6 +25,7 @@
* Fix ElasticSearch scroller bug.
* Add component ID for Aerospike(ID=149).
* Packages with name `recevier` are renamed to `receiver`.
* `BanyanDBMetricsDAO` handles `storeIDTag` in `multiGet` for `BanyanDBModelExtension`.
#### UI

View File

@ -43,6 +43,7 @@ import org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractB
import java.io.IOException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
@ -68,25 +69,44 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
if (schema == null) {
throw new IOException(model.getName() + " is not registered");
}
final Map<String, List<String>> seriesIDColumns = new HashMap<>();
if (model.getBanyanDBModelExtension().isStoreIDTag()) {
seriesIDColumns.put(BanyanDBConverter.ID, new ArrayList<>());
} else {
model.getColumns().forEach(c -> {
BanyanDBExtension ext = c.getBanyanDBExtension();
if (ext == null) {
return;
}
if (ext.isShardingKey()) {
seriesIDColumns.put(c.getColumnName().getName(), new ArrayList<>());
}
});
if (seriesIDColumns.isEmpty()) {
seriesIDColumns.put(ENTITY_ID, new ArrayList<>());
}
}
String tc = model.getBanyanDBModelExtension().getTimestampColumn();
final String tsCol = Strings.isBlank(tc) ? TIME_BUCKET : tc;
long begin = 0L, end = 0L;
final Map<String, List<String>> seriesIDColumns = new HashMap<>();
model.getColumns().forEach(c -> {
BanyanDBExtension ext = c.getBanyanDBExtension();
if (ext == null) {
return;
}
if (ext.isShardingKey()) {
seriesIDColumns.put(c.getColumnName().getName(), new ArrayList<>());
}
});
if (seriesIDColumns.isEmpty()) {
seriesIDColumns.put(ENTITY_ID, new ArrayList<>());
}
StringBuilder idStr = new StringBuilder();
for (Metrics m : metrics) {
AnalyticalResult result = analyze(m, tsCol, seriesIDColumns);
List<StorageID.Fragment> fragments = m.id().read();
if (model.getBanyanDBModelExtension().isStoreIDTag()) {
if (fragments.size() != 1) {
log.error("[{}]fragments' size is more than expected", fragments);
continue;
}
Object val = fragments.get(0).getValue();
fragments = Arrays.asList(new StorageID.Fragment(
new String[]{BanyanDBConverter.ID},
String.class,
true,
val));
}
AnalyticalResult result = analyze(fragments, tsCol, seriesIDColumns);
idStr.append(result.cols()).append("=").append(m.id().build()).append(",");
if (!result.success) {
continue;
@ -180,9 +200,7 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
}
}
private AnalyticalResult analyze(Metrics m, String tsCol, Map<String, List<String>> seriesIDColumns) {
StorageID id = m.id();
List<StorageID.Fragment> fragments = id.read();
private AnalyticalResult analyze(List<StorageID.Fragment> fragments, String tsCol, Map<String, List<String>> seriesIDColumns) {
AnalyticalResult result = new AnalyticalResult();
for (StorageID.Fragment f : fragments) {
Optional<String[]> cols = f.getName();
@ -202,12 +220,12 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
Preconditions.checkState(f.getType().equals(String.class));
seriesIDColumns.get(col).add((String) f.getValue());
} else {
log.error("col [{}] in fragment [{}] in id [{}] is not ts or seriesID", col, f, id.build());
log.error("col [{}] in fragment [{}] id [{}] is not ts or seriesID", col, f, fragments);
return result;
}
}
} else {
log.error("fragment [{}] in id [{}] doesn't contains cols", f, id.build());
log.error("fragment [{}] in id [{}] doesn't contains cols", f, fragments);
return result;
}
}