diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md
index eadee42fa8..6f8f53595a 100644
--- a/docs/en/changes/changes.md
+++ b/docs/en/changes/changes.md
@@ -85,6 +85,7 @@
* Skip persisting metrics/record data that have been expired.
* Fix the issue of missing Last Ping data.
* Add HTTP headers configuration for the alarm webhook.
+* Bump up BanyanDB java client to 0.8.0-rc3.
#### UI
diff --git a/oap-server-bom/pom.xml b/oap-server-bom/pom.xml
index 71b927ba1c..7ea6c6baca 100644
--- a/oap-server-bom/pom.xml
+++ b/oap-server-bom/pom.xml
@@ -72,7 +72,7 @@
4.4.16
4.1.5
1.21
- 0.8.0-rc2
+ 0.8.0-rc3
3.4.0
2.4.6.RELEASE
1.5.3
diff --git a/oap-server/server-starter/src/main/resources/bydb.yml b/oap-server/server-starter/src/main/resources/bydb.yml
index 3cfdb7da5a..6ecaa5a835 100644
--- a/oap-server/server-starter/src/main/resources/bydb.yml
+++ b/oap-server/server-starter/src/main/resources/bydb.yml
@@ -87,3 +87,7 @@ groups:
shardNum: ${SW_STORAGE_BANYANDB_GM_INDEX_SHARD_NUM:2}
segmentInterval: ${SW_STORAGE_BANYANDB_GM_INDEX_SI_DAYS:15}
ttl: ${SW_STORAGE_BANYANDB_GM_INDEX_TTL_DAYS:15}
+
+ # The group settings of UI and profiling.
+ property:
+ shardNum: ${SW_STORAGE_BANYANDB_GM_INDEX_SHARD_NUM:1}
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConfigLoader.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConfigLoader.java
index 34036ab4a6..12d5dcbbae 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConfigLoader.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConfigLoader.java
@@ -111,6 +111,10 @@ public class BanyanDBConfigLoader {
config.getMetadata(), (Properties) groups.get("metadata"),
moduleProvider.getModule().name(), moduleProvider.name()
);
+ copyProperties(
+ config.getProperty(), (Properties) groups.get("property"),
+ moduleProvider.getModule().name(), moduleProvider.name()
+ );
} catch (IllegalAccessException e) {
throw new ModuleStartException("Failed to load BanyanDB configuration.", e);
}
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
index 806903b99f..dd0d54d3fe 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBIndexInstaller.java
@@ -32,6 +32,7 @@ import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.Group;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.IntervalRule;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase;
+import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Property;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Measure;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Stream;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.IndexRule;
@@ -68,11 +69,6 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
public InstallInfo isExists(Model model) throws StorageException {
InstallInfoBanyanDB installInfo = new InstallInfoBanyanDB(model);
installInfo.setDownSampling(model.getDownsampling());
- if (!model.isTimeSeries()) {
- installInfo.setTableName(model.getName());
- installInfo.setAllExist(true);
- return installInfo;
- }
final DownSamplingConfigService downSamplingConfigService = moduleManager.find(CoreModule.NAME)
.provider()
.getService(DownSamplingConfigService.class);
@@ -91,35 +87,42 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
installInfo.setAllExist(false);
return installInfo;
} else {
- // register models only locally(Schema cache) but not remotely
- if (model.isRecord()) { // stream
- StreamModel streamModel = MetadataRegistry.INSTANCE.registerStreamModel(
- model, config, downSamplingConfigService);
- if (!RunningMode.isNoInitMode()) {
- checkStream(streamModel.getStream(), c);
- checkIndexRules(model.getName(), streamModel.getIndexRules(), c);
- checkIndexRuleBinding(
- streamModel.getIndexRules(), metadata.getGroup(), metadata.name(),
- BanyandbCommon.Catalog.CATALOG_STREAM, c
- );
- // Stream not support server side TopN pre-aggregation
+ if (model.isTimeSeries()) {
+ // register models only locally(Schema cache) but not remotely
+ if (model.isRecord()) { // stream
+ StreamModel streamModel = MetadataRegistry.INSTANCE.registerStreamModel(
+ model, config, downSamplingConfigService);
+ if (!RunningMode.isNoInitMode()) {
+ checkStream(streamModel.getStream(), c);
+ checkIndexRules(model.getName(), streamModel.getIndexRules(), c);
+ checkIndexRuleBinding(
+ streamModel.getIndexRules(), metadata.getGroup(), metadata.name(),
+ BanyandbCommon.Catalog.CATALOG_STREAM, c
+ );
+ // Stream not support server side TopN pre-aggregation
+ }
+ } else { // measure
+ MeasureModel measureModel = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, downSamplingConfigService);
+ if (!RunningMode.isNoInitMode()) {
+ checkMeasure(measureModel.getMeasure(), c);
+ checkIndexRules(model.getName(), measureModel.getIndexRules(), c);
+ checkIndexRuleBinding(
+ measureModel.getIndexRules(), metadata.getGroup(), metadata.name(),
+ BanyandbCommon.Catalog.CATALOG_MEASURE, c
+ );
+ checkTopNAggregation(model, c);
+ }
}
- } else { // measure
- MeasureModel measureModel = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, downSamplingConfigService);
- if (!RunningMode.isNoInitMode()) {
- checkMeasure(measureModel.getMeasure(), c);
- checkIndexRules(model.getName(), measureModel.getIndexRules(), c);
- checkIndexRuleBinding(
- measureModel.getIndexRules(), metadata.getGroup(), metadata.name(),
- BanyandbCommon.Catalog.CATALOG_MEASURE, c
- );
- checkTopNAggregation(model, c);
+ // pre-load remote schema for java client
+ MetadataCache.EntityMetadata remoteMeta = updateSchemaFromServer(metadata, c);
+ if (remoteMeta == null) {
+ throw new IllegalStateException("inconsistent state: metadata:" + metadata + ", remoteMeta: null");
+ }
+ } else {
+ PropertyModel propertyModel = MetadataRegistry.INSTANCE.registerPropertyModel(model, config);
+ if (!RunningMode.isNoInitMode()) {
+ checkProperty(propertyModel.getProperty(), c);
}
- }
- // pre-load remote schema for java client
- MetadataCache.EntityMetadata remoteMeta = updateSchemaFromServer(metadata, c);
- if (remoteMeta == null) {
- throw new IllegalStateException("inconsistent state: metadata:" + metadata + ", remoteMeta: null");
}
installInfo.setAllExist(true);
return installInfo;
@@ -135,73 +138,89 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
DownSamplingConfigService configService = moduleManager.find(CoreModule.NAME)
.provider()
.getService(DownSamplingConfigService.class);
- if (model.isRecord()) { // stream
- StreamModel streamModel = MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService);
- Stream stream = streamModel.getStream();
- if (stream != null) {
- log.info("install stream schema {}", model.getName());
- final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
- try {
- client.define(stream);
- if (CollectionUtils.isNotEmpty(streamModel.getIndexRules())) {
- for (IndexRule indexRule : streamModel.getIndexRules()) {
- defineIndexRule(model.getName(), indexRule, client);
+ if (model.isTimeSeries()) {
+ if (model.isRecord()) { // stream
+ StreamModel streamModel = MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService);
+ Stream stream = streamModel.getStream();
+ if (stream != null) {
+ log.info("install stream schema {}", model.getName());
+ final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
+ try {
+ client.define(stream);
+ if (CollectionUtils.isNotEmpty(streamModel.getIndexRules())) {
+ for (IndexRule indexRule : streamModel.getIndexRules()) {
+ defineIndexRule(model.getName(), indexRule, client);
+ }
+ defineIndexRuleBinding(
+ streamModel.getIndexRules(), stream.getMetadata().getGroup(), stream.getMetadata().getName(),
+ BanyandbCommon.Catalog.CATALOG_STREAM, client
+ );
+ }
+ } catch (BanyanDBException ex) {
+ if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
+ log.info(
+ "Stream schema {}_{} already created by another OAP node",
+ model.getName(),
+ model.getDownsampling()
+ );
+ } else {
+ throw ex;
}
- defineIndexRuleBinding(
- streamModel.getIndexRules(), stream.getMetadata().getGroup(), stream.getMetadata().getName(),
- BanyandbCommon.Catalog.CATALOG_STREAM, client
- );
}
- } catch (BanyanDBException ex) {
- if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
- log.info(
- "Stream schema {}_{} already created by another OAP node",
- model.getName(),
- model.getDownsampling()
- );
- } else {
- throw ex;
+ }
+ } else { // measure
+ MeasureModel measureModel = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService);
+ Measure measure = measureModel.getMeasure();
+ if (measure != null) {
+ log.info("install measure schema {}", model.getName());
+ final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
+ try {
+ client.define(measure);
+ if (CollectionUtils.isNotEmpty(measureModel.getIndexRules())) {
+ for (IndexRule indexRule : measureModel.getIndexRules()) {
+ defineIndexRule(model.getName(), indexRule, client);
+ }
+ defineIndexRuleBinding(
+ measureModel.getIndexRules(), measure.getMetadata().getGroup(), measure.getMetadata().getName(),
+ BanyandbCommon.Catalog.CATALOG_MEASURE, client
+ );
+ }
+ } catch (BanyanDBException ex) {
+ if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
+ log.info("Measure schema {}_{} already created by another OAP node",
+ model.getName(),
+ model.getDownsampling());
+ } else {
+ throw ex;
+ }
+ }
+ final MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(model);
+ try {
+ defineTopNAggregation(schema, client);
+ } catch (BanyanDBException ex) {
+ if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
+ log.info("Measure schema {}_{} TopN({}) already created by another OAP node",
+ model.getName(),
+ model.getDownsampling(),
+ schema.getTopNSpec());
+ } else {
+ throw ex;
+ }
}
}
}
- } else { // measure
- MeasureModel measureModel = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService);
- Measure measure = measureModel.getMeasure();
- if (measure != null) {
- log.info("install measure schema {}", model.getName());
- final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
- try {
- client.define(measure);
- if (CollectionUtils.isNotEmpty(measureModel.getIndexRules())) {
- for (IndexRule indexRule : measureModel.getIndexRules()) {
- defineIndexRule(model.getName(), indexRule, client);
- }
- defineIndexRuleBinding(
- measureModel.getIndexRules(), measure.getMetadata().getGroup(), measure.getMetadata().getName(),
- BanyandbCommon.Catalog.CATALOG_MEASURE, client
- );
- }
- } catch (BanyanDBException ex) {
- if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
- log.info("Measure schema {}_{} already created by another OAP node",
- model.getName(),
- model.getDownsampling());
- } else {
- throw ex;
- }
- }
- final MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(model);
- try {
- defineTopNAggregation(schema, client);
- } catch (BanyanDBException ex) {
- if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
- log.info("Measure schema {}_{} TopN({}) already created by another OAP node",
- model.getName(),
- model.getDownsampling(),
- schema.getTopNSpec());
- } else {
- throw ex;
- }
+ } else {
+ PropertyModel propertyModel = MetadataRegistry.INSTANCE.registerPropertyModel(model, config);
+ Property property = propertyModel.getProperty();
+ log.info("install property schema {}", model.getName());
+ final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
+ try {
+ client.define(property);
+ } catch (BanyanDBException ex) {
+ if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
+ log.info("Property schema {} already created by another OAP node", model.getName());
+ } else {
+ throw ex;
}
}
}
@@ -225,33 +244,50 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
ResourceExist resourceExist;
Group.Builder gBuilder
= Group.newBuilder()
- .setMetadata(BanyandbCommon.Metadata.newBuilder().setName(metadata.getGroup()))
- .setResourceOpts(BanyandbCommon.ResourceOpts.newBuilder()
- .setShardNum(metadata.getShard())
- .setSegmentInterval(
- IntervalRule.newBuilder()
- .setUnit(
- IntervalRule.Unit.UNIT_DAY)
- .setNum(
- metadata.getSegmentIntervalDays()))
- .setTtl(
- IntervalRule.newBuilder()
- .setUnit(
- IntervalRule.Unit.UNIT_DAY)
- .setNum(
- metadata.getTtlDays())));
+ .setMetadata(BanyandbCommon.Metadata.newBuilder().setName(metadata.getGroup()));
+ BanyandbCommon.ResourceOpts.Builder optsBuilder = BanyandbCommon.ResourceOpts.newBuilder().setShardNum(metadata.getShard());
+
switch (metadata.getKind()) {
case STREAM:
+ optsBuilder.setSegmentInterval(
+ IntervalRule.newBuilder()
+ .setUnit(
+ IntervalRule.Unit.UNIT_DAY)
+ .setNum(
+ metadata.getSegmentIntervalDays()))
+ .setTtl(
+ IntervalRule.newBuilder()
+ .setUnit(
+ IntervalRule.Unit.UNIT_DAY)
+ .setNum(
+ metadata.getTtlDays()));
resourceExist = client.existStream(metadata.getGroup(), metadata.name());
gBuilder.setCatalog(BanyandbCommon.Catalog.CATALOG_STREAM).build();
break;
case MEASURE:
+ optsBuilder.setSegmentInterval(
+ IntervalRule.newBuilder()
+ .setUnit(
+ IntervalRule.Unit.UNIT_DAY)
+ .setNum(
+ metadata.getSegmentIntervalDays()))
+ .setTtl(
+ IntervalRule.newBuilder()
+ .setUnit(
+ IntervalRule.Unit.UNIT_DAY)
+ .setNum(
+ metadata.getTtlDays()));
resourceExist = client.existMeasure(metadata.getGroup(), metadata.name());
gBuilder.setCatalog(BanyandbCommon.Catalog.CATALOG_MEASURE).build();
break;
+ case PROPERTY:
+ resourceExist = client.existProperty(metadata.getGroup(), metadata.name());
+ gBuilder.setCatalog(BanyandbCommon.Catalog.CATALOG_PROPERTY).build();
+ break;
default:
throw new IllegalStateException("unknown metadata kind: " + metadata.getKind());
}
+ gBuilder.setResourceOpts(optsBuilder.build());
if (!RunningMode.isNoInitMode()) {
if (!groupAligned.contains(metadata.getGroup())) {
// create the group if not exist
@@ -431,6 +467,26 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
}
}
+ /**
+ * Check if the property exists and update it if necessary
+ */
+ private void checkProperty(Property property, BanyanDBClient client) throws BanyanDBException {
+ Property hisProperty = client.findPropertyDefinition(property.getMetadata().getGroup(), property.getMetadata().getName());
+ if (hisProperty == null) {
+ throw new IllegalStateException("Property: " + property.getMetadata().getName() + " exist but can't find it from BanyanDB server");
+ } else {
+ boolean equals = hisProperty.toBuilder()
+ .clearUpdatedAt()
+ .clearMetadata()
+ .build()
+ .equals(property.toBuilder().clearUpdatedAt().clearMetadata().build());
+ if (!equals) {
+ client.update(property);
+ log.info("update Property: {} from: {} to: {}", hisProperty.getMetadata().getName(), hisProperty, property);
+ }
+ }
+ }
+
/**
* Check if the index rules exist and update them if necessary
*/
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageClient.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageClient.java
index b2fb795691..9bca7b35a1 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageClient.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageClient.java
@@ -128,7 +128,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
BanyandbProperty.QueryResponse resp
= this.client.query(BanyandbProperty.QueryRequest.newBuilder()
.addGroups(group)
- .setContainer(name)
+ .setName(name)
.setLimit(Integer.MAX_VALUE)
.build());
this.healthChecker.health();
@@ -148,7 +148,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
try {
BanyandbProperty.QueryResponse resp = this.client.query(BanyandbProperty.QueryRequest.newBuilder()
.addGroups(group)
- .setContainer(name)
+ .setName(name)
.addIds(id)
.build());
this.healthChecker.health();
@@ -225,7 +225,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
/**
* PropertyStore.Strategy is default to {@link Strategy#STRATEGY_MERGE}
*/
- public void define(Property property) throws IOException {
+ public void apply(Property property) throws IOException {
try {
this.client.apply(property);
this.healthChecker.health();
@@ -235,7 +235,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
}
}
- public void define(Property property, Strategy strategy) throws IOException {
+ public void apply(Property property, Strategy strategy) throws IOException {
try {
this.client.apply(property, strategy);
this.healthChecker.health();
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageConfig.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageConfig.java
index 373ff6f06b..66858e7aa6 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageConfig.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageConfig.java
@@ -27,6 +27,9 @@ import org.apache.skywalking.oap.server.library.module.ModuleConfig;
@Getter
@Setter
public class BanyanDBStorageConfig extends ModuleConfig {
+
+ public static final String PROPERTY_GROUP_NAME = "property";
+
private Global global = new Global();
private RecordsNormal recordsNormal = new RecordsNormal();
private RecordsSuper recordsSuper = new RecordsSuper();
@@ -34,6 +37,7 @@ public class BanyanDBStorageConfig extends ModuleConfig {
private MetricsHour metricsHour = new MetricsHour();
private MetricsDay metricsDay = new MetricsDay();
private Metadata metadata = new Metadata();
+ private Property property = new Property();
public String[] getTargetArray() {
return Iterables.toArray(
@@ -168,4 +172,13 @@ public class BanyanDBStorageConfig extends ModuleConfig {
private int segmentInterval = 15;
private int ttl = 15;
}
+
+ /**
+ * The group settings of UI and profiling.
+ */
+ @Getter
+ @Setter
+ public static class Property {
+ private int shardNum = 1;
+ }
}
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageProvider.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageProvider.java
index 5b6392e581..d224e7498a 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageProvider.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBStorageProvider.java
@@ -18,7 +18,6 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb;
-import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.storage.IBatchDAO;
import org.apache.skywalking.oap.server.core.storage.IHistoryDeleteDAO;
@@ -206,16 +205,6 @@ public class BanyanDBStorageProvider extends ModuleProvider {
this.client.registerChecker(healthChecker);
try {
this.client.connect();
- this.client.defineIfEmpty(BanyandbCommon.Group.newBuilder()
- .setMetadata(
- BanyandbCommon.Metadata.newBuilder()
- .setName(
- BanyanDBUITemplateManagementDAO.GROUP))
- .setCatalog(BanyandbCommon.Catalog.CATALOG_PROPERTY)
- .setResourceOpts(BanyandbCommon.ResourceOpts.newBuilder()
- .setShardNum(1)
- .build())
- .build());
this.modelInstaller.start();
getManager().find(CoreModule.NAME).provider().getService(ModelCreator.class).addModelListener(modelInstaller);
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUIMenuManagementDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUIMenuManagementDAO.java
index 74e7cc75dd..029c357fb4 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUIMenuManagementDAO.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUIMenuManagementDAO.java
@@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.banyandb.model.v1.BanyandbModel;
-import org.apache.skywalking.banyandb.property.v1.BanyandbProperty;
import org.apache.skywalking.banyandb.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.oap.server.core.management.ui.menu.UIMenu;
@@ -32,7 +31,6 @@ import java.io.IOException;
@Slf4j
public class BanyanDBUIMenuManagementDAO extends AbstractBanyanDBDAO implements UIMenuManagementDAO {
- public static final String GROUP = "sw";
public BanyanDBUIMenuManagementDAO(BanyanDBStorageClient client) {
super(client);
@@ -40,7 +38,7 @@ public class BanyanDBUIMenuManagementDAO extends AbstractBanyanDBDAO implements
@Override
public UIMenu getMenu(String id) throws IOException {
- Property p = getClient().queryProperty(GROUP, UIMenu.INDEX_NAME, id);
+ Property p = getClient().queryProperty(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, UIMenu.INDEX_NAME, id);
if (p == null) {
return null;
}
@@ -50,23 +48,19 @@ public class BanyanDBUIMenuManagementDAO extends AbstractBanyanDBDAO implements
@Override
public void saveMenu(UIMenu menu) throws IOException {
Property property = Property.newBuilder()
- .setMetadata(BanyandbProperty.Metadata.newBuilder().setId(menu.getMenuId())
- .setContainer(
- BanyandbCommon.Metadata.newBuilder()
- .setGroup(GROUP)
- .setName(
- UIMenu.INDEX_NAME)))
-
+ .setMetadata(
+ BanyandbCommon.Metadata.newBuilder().setGroup(BanyanDBStorageConfig.PROPERTY_GROUP_NAME).setName(UIMenu.INDEX_NAME))
+ .setId(menu.getMenuId())
.addTags(TagAndValue.newStringTag(UIMenu.CONFIGURATION, menu.getConfigurationJson())
.build())
.addTags(TagAndValue.newLongTag(UIMenu.UPDATE_TIME, menu.getUpdateTime()).build())
.build();
- this.getClient().define(property);
+ this.getClient().apply(property);
}
public UIMenu parse(Property property) {
UIMenu menu = new UIMenu();
- menu.setMenuId(property.getMetadata().getId());
+ menu.setMenuId(property.getId());
for (BanyandbModel.Tag tag : property.getTagsList()) {
TagAndValue> tagAndValue = TagAndValue.fromProtobuf(tag);
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUITemplateManagementDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUITemplateManagementDAO.java
index 9bd4ef26bf..71c121d178 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUITemplateManagementDAO.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBUITemplateManagementDAO.java
@@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.banyandb.model.v1.BanyandbModel;
-import org.apache.skywalking.banyandb.property.v1.BanyandbProperty;
import org.apache.skywalking.banyandb.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.oap.server.core.management.ui.template.UITemplate;
@@ -37,7 +36,6 @@ import java.util.stream.Collectors;
@Slf4j
public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO implements UITemplateManagementDAO {
- public static final String GROUP = "sw";
public BanyanDBUITemplateManagementDAO(BanyanDBStorageClient client) {
super(client);
@@ -45,7 +43,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
@Override
public DashboardConfiguration getTemplate(String id) throws IOException {
- Property p = getClient().queryProperty(GROUP, UITemplate.INDEX_NAME, id);
+ Property p = getClient().queryProperty(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, UITemplate.INDEX_NAME, id);
if (p == null) {
return null;
}
@@ -54,7 +52,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
@Override
public List getAllTemplates(Boolean includingDisabled) throws IOException {
- List propertyList = getClient().listProperties(GROUP, UITemplate.INDEX_NAME);
+ List propertyList = getClient().listProperties(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, UITemplate.INDEX_NAME);
return propertyList.stream().map(p -> fromEntity(parse(p)))
.filter(conf -> includingDisabled || !conf.isDisabled())
.collect(Collectors.toList());
@@ -64,10 +62,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
public TemplateChangeStatus addTemplate(DashboardSetting setting) {
Property newTemplate = applyAll(setting.toEntity());
try {
- this.getClient().define(newTemplate);
+ this.getClient().apply(newTemplate);
return TemplateChangeStatus.builder()
.status(true)
- .id(newTemplate.getMetadata().getId())
+ .id(newTemplate.getId())
.build();
} catch (IOException ioEx) {
log.error("fail to add new template", ioEx);
@@ -80,10 +78,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
public TemplateChangeStatus changeTemplate(DashboardSetting setting) {
Property newTemplate = applyConfiguration(setting.toEntity());
try {
- this.getClient().define(newTemplate);
+ this.getClient().apply(newTemplate);
return TemplateChangeStatus.builder()
.status(true)
- .id(newTemplate.getMetadata().getId())
+ .id(newTemplate.getId())
.build();
} catch (IOException ioEx) {
log.error("fail to modify the template", ioEx);
@@ -94,14 +92,14 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
@Override
public TemplateChangeStatus disableTemplate(String id) throws IOException {
- Property oldProperty = this.getClient().queryProperty(GROUP, UITemplate.INDEX_NAME, id);
+ Property oldProperty = this.getClient().queryProperty(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, UITemplate.INDEX_NAME, id);
if (oldProperty == null) {
return TemplateChangeStatus.builder().status(false).id(id).message("Can't find the template")
.build();
}
UITemplate uiTemplate = parse(oldProperty);
try {
- this.getClient().define(applyStatus(uiTemplate));
+ this.getClient().apply(applyStatus(uiTemplate));
return TemplateChangeStatus.builder()
.status(true)
.id(uiTemplate.id().build())
@@ -121,7 +119,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
public UITemplate parse(Property property) {
UITemplate uiTemplate = new UITemplate();
- uiTemplate.setTemplateId(property.getMetadata().getId());
+ uiTemplate.setTemplateId(property.getId());
for (BanyandbModel.Tag tag : property.getTagsList()) {
TagAndValue> tagAndValue = TagAndValue.fromProtobuf(tag);
@@ -138,11 +136,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
public Property applyAll(UITemplate uiTemplate) {
return Property.newBuilder()
- .setMetadata(BanyandbProperty.Metadata.newBuilder()
- .setId(uiTemplate.id().build())
- .setContainer(BanyandbCommon.Metadata.newBuilder()
- .setGroup(GROUP)
- .setName(UITemplate.INDEX_NAME)))
+ .setMetadata(BanyandbCommon.Metadata.newBuilder()
+ .setGroup(BanyanDBStorageConfig.PROPERTY_GROUP_NAME)
+ .setName(UITemplate.INDEX_NAME))
+ .setId(uiTemplate.id().build())
.addTags(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration()).build())
.addTags(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled()).build())
.addTags(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()).build())
@@ -157,11 +154,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
*/
public Property applyStatus(UITemplate uiTemplate) {
return Property.newBuilder()
- .setMetadata(BanyandbProperty.Metadata.newBuilder()
- .setId(uiTemplate.id().build())
- .setContainer(BanyandbCommon.Metadata.newBuilder()
- .setGroup(GROUP)
- .setName(UITemplate.INDEX_NAME)))
+ .setMetadata(BanyandbCommon.Metadata.newBuilder()
+ .setGroup(BanyanDBStorageConfig.PROPERTY_GROUP_NAME)
+ .setName(UITemplate.INDEX_NAME))
+ .setId(uiTemplate.id().build())
.addTags(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled()).build())
.addTags(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()).build())
.build();
@@ -175,11 +171,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
*/
public Property applyConfiguration(UITemplate uiTemplate) {
return Property.newBuilder()
- .setMetadata(BanyandbProperty.Metadata.newBuilder()
- .setId(uiTemplate.id().build())
- .setContainer(BanyandbCommon.Metadata.newBuilder()
- .setGroup(GROUP)
- .setName(UITemplate.INDEX_NAME)))
+ .setMetadata(BanyandbCommon.Metadata.newBuilder()
+ .setGroup(BanyanDBStorageConfig.PROPERTY_GROUP_NAME)
+ .setName(UITemplate.INDEX_NAME))
+ .setId(uiTemplate.id().build())
.addTags(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration()).build())
.addTags(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()).build())
.build();
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
index f88c7619ea..882c9cfef2 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java
@@ -30,6 +30,7 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.Metadata;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase;
+import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Property;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.CompressionMethod;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.EncodingMethod;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.FieldSpec;
@@ -63,6 +64,7 @@ import javax.annotation.Nonnull;
import javax.annotation.Nullable;
import java.lang.reflect.ParameterizedType;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.Comparator;
import java.util.HashMap;
import java.util.List;
@@ -177,6 +179,19 @@ public enum MetadataRegistry {
return new MeasureModel(builder.build(), indexRules);
}
+ public PropertyModel registerPropertyModel(Model model, BanyanDBStorageConfig config) {
+ final SchemaMetadata schemaMetadata = parseMetadata(model, config, null);
+ Schema.SchemaBuilder schemaBuilder = Schema.builder().metadata(schemaMetadata);
+ List tags = parseTagMetadata(model, schemaBuilder, Collections.emptyList(), schemaMetadata.group);
+ final Property.Builder builder = Property.newBuilder();
+ builder.setMetadata(BanyandbCommon.Metadata.newBuilder().setGroup(schemaMetadata.getGroup())
+ .setName(schemaMetadata.name()));
+ for (TagMetadata tag : tags) {
+ builder.addTags(tag.getTagSpec());
+ }
+ return new PropertyModel(builder.build());
+ }
+
private TopNAggregation parseTopNSpec(final Model model, final String group, final String measureName)
throws StorageException {
if (model.getBanyanDBModelExtension().getTopN() == null) {
@@ -434,7 +449,7 @@ public enum MetadataRegistry {
TagSpec.Builder tagSpec = TagSpec.newBuilder().setName(colName);
if (String.class.equals(clazz) || StorageDataComplexObject.class.isAssignableFrom(clazz) || JsonObject.class.equals(clazz)) {
tagSpec = tagSpec.setType(TagType.TAG_TYPE_STRING);
- } else if (int.class.equals(clazz) || long.class.equals(clazz)) {
+ } else if (int.class.equals(clazz) || long.class.equals(clazz) || Integer.class.equals(clazz) || Long.class.equals(clazz)) {
tagSpec = tagSpec.setType(TagType.TAG_TYPE_INT);
} else if (byte[].class.equals(clazz)) {
tagSpec = tagSpec.setType(TagType.TAG_TYPE_DATA_BINARY);
@@ -461,6 +476,9 @@ public enum MetadataRegistry {
}
public SchemaMetadata parseMetadata(Model model, BanyanDBStorageConfig config, DownSamplingConfigService configService) {
+ if (!model.isTimeSeries()) {
+ return new SchemaMetadata(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, model.getName(), Kind.PROPERTY, DownSampling.None, config.getProperty().getShardNum(), 0, 0);
+ }
if (model.isRecord()) { // stream
return new SchemaMetadata(model.isSuperDataset() ? model.getName() : "normal",
model.getName(),
@@ -600,7 +618,7 @@ public enum MetadataRegistry {
}
public enum Kind {
- MEASURE, STREAM;
+ MEASURE, STREAM, PROPERTY;
}
@RequiredArgsConstructor
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/PropertyModel.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/PropertyModel.java
new file mode 100644
index 0000000000..073066c224
--- /dev/null
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/PropertyModel.java
@@ -0,0 +1,30 @@
+/*
+ * 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.skywalking.oap.server.storage.plugin.banyandb;
+
+import lombok.Getter;
+import lombok.RequiredArgsConstructor;
+import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Property;
+
+@RequiredArgsConstructor
+@Getter
+public class PropertyModel {
+ private final Property property;
+}
diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBContinuousProfilingPolicyDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBContinuousProfilingPolicyDAO.java
index 52e69a4574..9cea876fbb 100644
--- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBContinuousProfilingPolicyDAO.java
+++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBContinuousProfilingPolicyDAO.java
@@ -21,12 +21,12 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.banyandb.model.v1.BanyandbModel;
-import org.apache.skywalking.banyandb.property.v1.BanyandbProperty;
import org.apache.skywalking.banyandb.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.oap.server.core.profiling.continuous.storage.ContinuousProfilingPolicy;
import org.apache.skywalking.oap.server.core.storage.profiling.continuous.IContinuousProfilingPolicyDAO;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
+import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageConfig;
import java.io.IOException;
import java.util.List;
@@ -35,7 +35,6 @@ import java.util.stream.Collectors;
@Slf4j
public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO implements IContinuousProfilingPolicyDAO {
- private static final String GROUP = "sw";
public BanyanDBContinuousProfilingPolicyDAO(BanyanDBStorageClient client) {
super(client);
@@ -44,7 +43,7 @@ public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO im
@Override
public void savePolicy(ContinuousProfilingPolicy policy) throws IOException {
try {
- this.getClient().define(applyAll(policy));
+ this.getClient().apply(applyAll(policy));
} catch (IOException e) {
log.error("fail to save policy", e);
}
@@ -52,11 +51,10 @@ public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO im
public Property applyAll(ContinuousProfilingPolicy policy) {
return Property.newBuilder()
- .setMetadata(BanyandbProperty.Metadata.newBuilder()
- .setId(policy.id().build())
- .setContainer(BanyandbCommon.Metadata.newBuilder()
- .setGroup(GROUP)
- .setName(ContinuousProfilingPolicy.INDEX_NAME)))
+ .setMetadata(BanyandbCommon.Metadata.newBuilder()
+ .setGroup(BanyanDBStorageConfig.PROPERTY_GROUP_NAME)
+ .setName(ContinuousProfilingPolicy.INDEX_NAME))
+ .setId(policy.id().build())
.addTags(TagAndValue.newStringTag(ContinuousProfilingPolicy.UUID, policy.getUuid()).build())
.addTags(TagAndValue.newStringTag(ContinuousProfilingPolicy.CONFIGURATION_JSON, policy.getConfigurationJson()).build())
.build();
@@ -66,14 +64,14 @@ public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO im
public List queryPolicies(List serviceIdList) throws IOException {
return serviceIdList.stream().map(s -> {
try {
- return getClient().queryProperty(GROUP, ContinuousProfilingPolicy.INDEX_NAME, s);
+ return getClient().queryProperty(BanyanDBStorageConfig.PROPERTY_GROUP_NAME, ContinuousProfilingPolicy.INDEX_NAME, s);
} catch (IOException e) {
log.warn("query policy error", e);
return null;
}
}).filter(Objects::nonNull).map(properties -> {
final ContinuousProfilingPolicy policy = new ContinuousProfilingPolicy();
- policy.setServiceId(properties.getMetadata().getId());
+ policy.setServiceId(properties.getId());
for (BanyandbModel.Tag tag : properties.getTagsList()) {
TagAndValue> tagAndValue = TagAndValue.fromProtobuf(tag);
if (tagAndValue.getTagName().equals(ContinuousProfilingPolicy.CONFIGURATION_JSON)) {
diff --git a/test/e2e-v2/script/env b/test/e2e-v2/script/env
index 76f6c1dbfe..486ad4ab31 100644
--- a/test/e2e-v2/script/env
+++ b/test/e2e-v2/script/env
@@ -23,7 +23,7 @@ SW_AGENT_CLIENT_JS_COMMIT=af0565a67d382b683c1dbd94c379b7080db61449
SW_AGENT_CLIENT_JS_TEST_COMMIT=4f1eb1dcdbde3ec4a38534bf01dded4ab5d2f016
SW_KUBERNETES_COMMIT_SHA=6fe5e6f0d3b7686c6be0457733e825ee68cb9b35
SW_ROVER_COMMIT=4c0cb8429a96f190ea30eac1807008d523c749c3
-SW_BANYANDB_COMMIT=f484391b33674fd5b6c6c1903056803818f5393a
+SW_BANYANDB_COMMIT=53b3be42d162e2f4ef0c667dc30f25e42ff17d70
SW_AGENT_PHP_COMMIT=3192c553002707d344bd6774cfab5bc61f67a1d3
SW_PREDICTOR_COMMIT=54a0197654a3781a6f73ce35146c712af297c994