Adapt BanyanDB Java Client 0.7.0. (#12621)

This commit is contained in:
Wan Kai 2024-09-14 15:05:20 +08:00 committed by GitHub
parent f8716b49e6
commit ddbed6d091
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
13 changed files with 389 additions and 172 deletions

View File

@ -65,6 +65,7 @@
* Fix the previous analysis result missing in the ALS `k8s-mesh` analyzer. * Fix the previous analysis result missing in the ALS `k8s-mesh` analyzer.
* Fix `findEndpoint` query require `keyword` when using BanyanDB. * Fix `findEndpoint` query require `keyword` when using BanyanDB.
* Support to analysis the ztunnel mapped IP address in eBPF Access Log Receiver. * Support to analysis the ztunnel mapped IP address in eBPF Access Log Receiver.
* Adapt BanyanDB Java Client 0.7.0-rc3.
#### UI #### UI

View File

@ -73,7 +73,7 @@
<httpcore.version>4.4.13</httpcore.version> <httpcore.version>4.4.13</httpcore.version>
<httpasyncclient.version>4.1.5</httpasyncclient.version> <httpasyncclient.version>4.1.5</httpasyncclient.version>
<commons-compress.version>1.21</commons-compress.version> <commons-compress.version>1.21</commons-compress.version>
<banyandb-java-client.version>0.7.0-rc2</banyandb-java-client.version> <banyandb-java-client.version>0.7-rc3</banyandb-java-client.version>
<kafka-clients.version>3.4.0</kafka-clients.version> <kafka-clients.version>3.4.0</kafka-clients.version>
<spring-kafka-test.version>2.4.6.RELEASE</spring-kafka-test.version> <spring-kafka-test.version>2.4.6.RELEASE</spring-kafka-test.version>
<consul.client.version>1.5.3</consul.client.version> <consul.client.version>1.5.3</consul.client.version>

View File

@ -22,8 +22,9 @@ import io.grpc.Status;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.v1.client.BanyanDBClient; import org.apache.skywalking.banyandb.v1.client.BanyanDBClient;
import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException; import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException;
import org.apache.skywalking.banyandb.v1.client.metadata.Measure; import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Measure;
import org.apache.skywalking.banyandb.v1.client.metadata.Stream; import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Stream;
import org.apache.skywalking.banyandb.v1.client.metadata.MetadataCache;
import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.config.ConfigService; import org.apache.skywalking.oap.server.core.config.ConfigService;
import org.apache.skywalking.oap.server.core.storage.StorageException; import org.apache.skywalking.oap.server.core.storage.StorageException;
@ -31,6 +32,7 @@ import org.apache.skywalking.oap.server.core.storage.model.Model;
import org.apache.skywalking.oap.server.core.storage.model.ModelInstaller; import org.apache.skywalking.oap.server.core.storage.model.ModelInstaller;
import org.apache.skywalking.oap.server.library.client.Client; import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
@Slf4j @Slf4j
public class BanyanDBIndexInstaller extends ModelInstaller { public class BanyanDBIndexInstaller extends ModelInstaller {
@ -58,20 +60,20 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
final boolean resourceExist = metadata.checkResourceExistence(c); final boolean resourceExist = metadata.checkResourceExistence(c);
if (!resourceExist) { if (!resourceExist) {
return false; return false;
} } else {
// register models only locally(Schema cache) but not remotely
// then check entity schema
if (metadata.findRemoteSchema(c).isPresent()) {
// register models only locally but not remotely
if (model.isRecord()) { // stream if (model.isRecord()) { // stream
MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService); MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService);
} else { // measure } else { // measure
MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService); MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService);
} }
// pre-load remote schema for java client
MetadataCache.EntityMetadata remoteMeta = metadata.updateRemoteSchema(c);
if (remoteMeta == null) {
throw new IllegalStateException("inconsistent state: metadata:" + metadata + ", remoteMeta: null");
}
return true; return true;
} }
throw new IllegalStateException("inconsistent state:" + metadata);
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
throw new StorageException("fail to check existence", ex); throw new StorageException("fail to check existence", ex);
} }
@ -84,11 +86,17 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
.provider() .provider()
.getService(ConfigService.class); .getService(ConfigService.class);
if (model.isRecord()) { // stream if (model.isRecord()) { // stream
Stream stream = MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService); StreamModel streamModel = MetadataRegistry.INSTANCE.registerStreamModel(model, config, configService);
Stream stream = streamModel.getStream();
if (stream != null) { if (stream != null) {
log.info("install stream schema {}", model.getName()); log.info("install stream schema {}", model.getName());
final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
try { try {
((BanyanDBStorageClient) client).define(stream); if (CollectionUtils.isNotEmpty(streamModel.getIndexRules())) {
client.define(stream, streamModel.getIndexRules());
} else {
client.define(stream);
}
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) { if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
log.info( log.info(
@ -102,12 +110,17 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
} }
} }
} else { // measure } else { // measure
Measure measure = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService); MeasureModel measureModel = MetadataRegistry.INSTANCE.registerMeasureModel(model, config, configService);
Measure measure = measureModel.getMeasure();
if (measure != null) { if (measure != null) {
log.info("install measure schema {}", measure.name()); log.info("install measure schema {}", model.getName());
final BanyanDBClient c = ((BanyanDBStorageClient) this.client).client; final BanyanDBClient client = ((BanyanDBStorageClient) this.client).client;
try { try {
c.define(measure); if (CollectionUtils.isNotEmpty(measureModel.getIndexRules())) {
client.define(measure, measureModel.getIndexRules());
} else {
client.define(measure);
}
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) { if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
log.info("Measure schema {}_{} already created by another OAP node", log.info("Measure schema {}_{} already created by another OAP node",
@ -119,7 +132,7 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
} }
final MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(model); final MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(model);
try { try {
schema.installTopNAggregation(c); schema.installTopNAggregation(client);
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) { if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
log.info("Measure schema {}_{} TopN({}) already created by another OAP node", log.info("Measure schema {}_{} TopN({}) already created by another OAP node",

View File

@ -44,7 +44,7 @@ public class BanyanDBNoneStreamDAO extends AbstractDAO<BanyanDBStorageClient> im
if (schema == null) { if (schema == null) {
throw new IOException(model.getName() + " is not registered"); throw new IOException(model.getName() + " is not registered");
} }
StreamWrite streamWrite = getClient().client.createStreamWrite( StreamWrite streamWrite = getClient().createStreamWrite(
schema.getMetadata().getGroup(), // group name schema.getMetadata().getGroup(), // group name
schema.getMetadata().name(), // stream-name schema.getMetadata().name(), // stream-name
noneStream.id().build() // identity noneStream.id().build() // identity

View File

@ -19,6 +19,7 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb; package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import io.grpc.Status; import io.grpc.Status;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase;
import org.apache.skywalking.banyandb.v1.client.BanyanDBClient; import org.apache.skywalking.banyandb.v1.client.BanyanDBClient;
import org.apache.skywalking.banyandb.v1.client.MeasureBulkWriteProcessor; import org.apache.skywalking.banyandb.v1.client.MeasureBulkWriteProcessor;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery; import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
@ -30,14 +31,15 @@ import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.StreamWrite; import org.apache.skywalking.banyandb.v1.client.StreamWrite;
import org.apache.skywalking.banyandb.v1.client.TopNQuery; import org.apache.skywalking.banyandb.v1.client.TopNQuery;
import org.apache.skywalking.banyandb.v1.client.TopNQueryResponse; import org.apache.skywalking.banyandb.v1.client.TopNQueryResponse;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.Group;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.TopNAggregation;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Measure;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Stream;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.ApplyRequest.Strategy;
import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.DeleteResponse;
import org.apache.skywalking.banyandb.v1.client.grpc.exception.AlreadyExistsException; import org.apache.skywalking.banyandb.v1.client.grpc.exception.AlreadyExistsException;
import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException; import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException;
import org.apache.skywalking.banyandb.v1.client.metadata.Group;
import org.apache.skywalking.banyandb.v1.client.metadata.Measure;
import org.apache.skywalking.banyandb.v1.client.metadata.Property;
import org.apache.skywalking.banyandb.v1.client.metadata.PropertyStore;
import org.apache.skywalking.banyandb.v1.client.metadata.Stream;
import org.apache.skywalking.banyandb.v1.client.metadata.TopNAggregation;
import org.apache.skywalking.oap.server.library.client.Client; import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.client.healthcheck.DelegatedHealthChecker; import org.apache.skywalking.oap.server.library.client.healthcheck.DelegatedHealthChecker;
import org.apache.skywalking.oap.server.library.client.healthcheck.HealthCheckable; import org.apache.skywalking.oap.server.library.client.healthcheck.HealthCheckable;
@ -103,9 +105,9 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
} }
public PropertyStore.DeleteResult deleteProperty(String group, String name, String id, String... tags) throws IOException { public DeleteResponse deleteProperty(String group, String name, String id, String... tags) throws IOException {
try { try {
PropertyStore.DeleteResult result = this.client.deleteProperty(group, name, id, tags); DeleteResponse result = this.client.deleteProperty(group, name, id, tags);
this.healthChecker.health(); this.healthChecker.health();
return result; return result;
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
@ -158,7 +160,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
/** /**
* PropertyStore.Strategy is default to {@link PropertyStore.Strategy#MERGE} * PropertyStore.Strategy is default to {@link Strategy#STRATEGY_MERGE}
*/ */
public void define(Property property) throws IOException { public void define(Property property) throws IOException {
try { try {
@ -170,7 +172,7 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
} }
public void define(Property property, PropertyStore.Strategy strategy) throws IOException { public void define(Property property, Strategy strategy) throws IOException {
try { try {
this.client.apply(property, strategy); this.client.apply(property, strategy);
this.healthChecker.health(); this.healthChecker.health();
@ -190,6 +192,16 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
} }
public void define(Stream stream, List<BanyandbDatabase.IndexRule> indexRules) throws BanyanDBException {
try {
this.client.define(stream, indexRules);
this.healthChecker.health();
} catch (BanyanDBException ex) {
healthChecker.unHealth(ex);
throw ex;
}
}
public void define(Measure measure) throws BanyanDBException { public void define(Measure measure) throws BanyanDBException {
try { try {
this.client.define(measure); this.client.define(measure);
@ -200,6 +212,16 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
} }
public void define(Measure measure, List<BanyandbDatabase.IndexRule> indexRules) throws BanyanDBException {
try {
this.client.define(measure, indexRules);
this.healthChecker.health();
} catch (BanyanDBException ex) {
healthChecker.unHealth(ex);
throw ex;
}
}
public void defineIfEmpty(Group group) throws IOException { public void defineIfEmpty(Group group) throws IOException {
try { try {
try { try {
@ -223,12 +245,20 @@ public class BanyanDBStorageClient implements Client, HealthCheckable {
} }
} }
public StreamWrite createStreamWrite(String group, String name, String elementId) { public StreamWrite createStreamWrite(String group, String name, String elementId) throws IOException {
return this.client.createStreamWrite(group, name, elementId); try {
return this.client.createStreamWrite(group, name, elementId);
} catch (BanyanDBException e) {
throw new IOException("fail to create stream write", e);
}
} }
public MeasureWrite createMeasureWrite(String group, String name, long timestamp) { public MeasureWrite createMeasureWrite(String group, String name, long timestamp) throws IOException {
return this.client.createMeasureWrite(group, name, timestamp); try {
return this.client.createMeasureWrite(group, name, timestamp);
} catch (BanyanDBException e) {
throw new IOException("fail to create measure write", e);
}
} }
public void write(StreamWrite streamWrite) { public void write(StreamWrite streamWrite) {

View File

@ -18,7 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb; package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import org.apache.skywalking.banyandb.v1.client.metadata.Group; import org.apache.skywalking.banyandb.common.v1.BanyandbCommon;
import org.apache.skywalking.oap.server.core.CoreModule; 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.IBatchDAO;
import org.apache.skywalking.oap.server.core.storage.IHistoryDeleteDAO; import org.apache.skywalking.oap.server.core.storage.IHistoryDeleteDAO;
@ -175,7 +175,13 @@ public class BanyanDBStorageProvider extends ModuleProvider {
this.client.registerChecker(healthChecker); this.client.registerChecker(healthChecker);
try { try {
this.client.connect(); this.client.connect();
this.client.defineIfEmpty(Group.create(BanyanDBUITemplateManagementDAO.GROUP)); this.client.defineIfEmpty(BanyandbCommon.Group.newBuilder()
.setMetadata(
BanyandbCommon.Metadata.newBuilder()
.setName(
BanyanDBUITemplateManagementDAO.GROUP))
.setCatalog(BanyandbCommon.Catalog.CATALOG_UNSPECIFIED)
.build());
this.modelInstaller.start(); this.modelInstaller.start();
getManager().find(CoreModule.NAME).provider().getService(ModelCreator.class).addModelListener(modelInstaller); getManager().find(CoreModule.NAME).provider().getService(ModelCreator.class).addModelListener(modelInstaller);

View File

@ -19,8 +19,11 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb; package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import lombok.extern.slf4j.Slf4j; 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.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.v1.client.metadata.Property; import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.oap.server.core.management.ui.menu.UIMenu; import org.apache.skywalking.oap.server.core.management.ui.menu.UIMenu;
import org.apache.skywalking.oap.server.core.storage.management.UIMenuManagementDAO; import org.apache.skywalking.oap.server.core.storage.management.UIMenuManagementDAO;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractBanyanDBDAO; import org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractBanyanDBDAO;
@ -46,17 +49,27 @@ public class BanyanDBUIMenuManagementDAO extends AbstractBanyanDBDAO implements
@Override @Override
public void saveMenu(UIMenu menu) throws IOException { public void saveMenu(UIMenu menu) throws IOException {
this.getClient().define(Property.create(GROUP, UIMenu.INDEX_NAME, menu.id().build()) Property property = Property.newBuilder()
.addTag(TagAndValue.newStringTag(UIMenu.CONFIGURATION, menu.getConfigurationJson())) .setMetadata(BanyandbProperty.Metadata.newBuilder().setId(menu.getMenuId())
.addTag(TagAndValue.newLongTag(UIMenu.UPDATE_TIME, menu.getUpdateTime())) .setContainer(
.build()); BanyandbCommon.Metadata.newBuilder()
.setGroup(GROUP)
.setName(
UIMenu.INDEX_NAME)))
.addTags(TagAndValue.newStringTag(UIMenu.CONFIGURATION, menu.getConfigurationJson())
.build())
.addTags(TagAndValue.newLongTag(UIMenu.UPDATE_TIME, menu.getUpdateTime()).build())
.build();
this.getClient().define(property);
} }
public UIMenu parse(Property property) { public UIMenu parse(Property property) {
UIMenu menu = new UIMenu(); UIMenu menu = new UIMenu();
menu.setMenuId(property.id()); menu.setMenuId(property.getMetadata().getId());
for (TagAndValue<?> tagAndValue : property.tags()) { for (BanyandbModel.Tag tag : property.getTagsList()) {
TagAndValue<?> tagAndValue = TagAndValue.fromProtobuf(tag);
if (tagAndValue.getTagName().equals(UIMenu.CONFIGURATION)) { if (tagAndValue.getTagName().equals(UIMenu.CONFIGURATION)) {
menu.setConfigurationJson((String) tagAndValue.getValue()); menu.setConfigurationJson((String) tagAndValue.getValue());
} else if (tagAndValue.getTagName().equals(UIMenu.UPDATE_TIME)) { } else if (tagAndValue.getTagName().equals(UIMenu.UPDATE_TIME)) {

View File

@ -19,8 +19,11 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb; package org.apache.skywalking.oap.server.storage.plugin.banyandb;
import lombok.extern.slf4j.Slf4j; 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.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.v1.client.metadata.Property; import org.apache.skywalking.banyandb.property.v1.BanyandbProperty.Property;
import org.apache.skywalking.oap.server.core.management.ui.template.UITemplate; import org.apache.skywalking.oap.server.core.management.ui.template.UITemplate;
import org.apache.skywalking.oap.server.core.query.input.DashboardSetting; import org.apache.skywalking.oap.server.core.query.input.DashboardSetting;
import org.apache.skywalking.oap.server.core.query.type.DashboardConfiguration; import org.apache.skywalking.oap.server.core.query.type.DashboardConfiguration;
@ -64,7 +67,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
this.getClient().define(newTemplate); this.getClient().define(newTemplate);
return TemplateChangeStatus.builder() return TemplateChangeStatus.builder()
.status(true) .status(true)
.id(newTemplate.id()) .id(newTemplate.getMetadata().getId())
.build(); .build();
} catch (IOException ioEx) { } catch (IOException ioEx) {
log.error("fail to add new template", ioEx); log.error("fail to add new template", ioEx);
@ -80,7 +83,7 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
this.getClient().define(newTemplate); this.getClient().define(newTemplate);
return TemplateChangeStatus.builder() return TemplateChangeStatus.builder()
.status(true) .status(true)
.id(newTemplate.id()) .id(newTemplate.getMetadata().getId())
.build(); .build();
} catch (IOException ioEx) { } catch (IOException ioEx) {
log.error("fail to modify the template", ioEx); log.error("fail to modify the template", ioEx);
@ -118,9 +121,10 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
public UITemplate parse(Property property) { public UITemplate parse(Property property) {
UITemplate uiTemplate = new UITemplate(); UITemplate uiTemplate = new UITemplate();
uiTemplate.setTemplateId(property.id()); uiTemplate.setTemplateId(property.getMetadata().getId());
for (TagAndValue<?> tagAndValue : property.tags()) { for (BanyandbModel.Tag tag : property.getTagsList()) {
TagAndValue<?> tagAndValue = TagAndValue.fromProtobuf(tag);
if (tagAndValue.getTagName().equals(UITemplate.CONFIGURATION)) { if (tagAndValue.getTagName().equals(UITemplate.CONFIGURATION)) {
uiTemplate.setConfiguration((String) tagAndValue.getValue()); uiTemplate.setConfiguration((String) tagAndValue.getValue());
} else if (tagAndValue.getTagName().equals(UITemplate.DISABLED)) { } else if (tagAndValue.getTagName().equals(UITemplate.DISABLED)) {
@ -133,11 +137,16 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
} }
public Property applyAll(UITemplate uiTemplate) { public Property applyAll(UITemplate uiTemplate) {
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build()) return Property.newBuilder()
.addTag(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration())) .setMetadata(BanyandbProperty.Metadata.newBuilder()
.addTag(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled())) .setId(uiTemplate.id().build())
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime())) .setContainer(BanyandbCommon.Metadata.newBuilder()
.build(); .setGroup(GROUP)
.setName(UITemplate.INDEX_NAME)))
.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())
.build();
} }
/** /**
@ -147,10 +156,15 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
* @return new property (patch) to be applied * @return new property (patch) to be applied
*/ */
public Property applyStatus(UITemplate uiTemplate) { public Property applyStatus(UITemplate uiTemplate) {
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build()) return Property.newBuilder()
.addTag(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled())) .setMetadata(BanyandbProperty.Metadata.newBuilder()
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime())) .setId(uiTemplate.id().build())
.build(); .setContainer(BanyandbCommon.Metadata.newBuilder()
.setGroup(GROUP)
.setName(UITemplate.INDEX_NAME)))
.addTags(TagAndValue.newLongTag(UITemplate.DISABLED, uiTemplate.getDisabled()).build())
.addTags(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()).build())
.build();
} }
/** /**
@ -160,9 +174,14 @@ public class BanyanDBUITemplateManagementDAO extends AbstractBanyanDBDAO impleme
* @return new property (patch) to be applied * @return new property (patch) to be applied
*/ */
public Property applyConfiguration(UITemplate uiTemplate) { public Property applyConfiguration(UITemplate uiTemplate) {
return Property.create(GROUP, UITemplate.INDEX_NAME, uiTemplate.id().build()) return Property.newBuilder()
.addTag(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration())) .setMetadata(BanyandbProperty.Metadata.newBuilder()
.addTag(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime())) .setId(uiTemplate.id().build())
.build(); .setContainer(BanyandbCommon.Metadata.newBuilder()
.setGroup(GROUP)
.setName(UITemplate.INDEX_NAME)))
.addTags(TagAndValue.newStringTag(UITemplate.CONFIGURATION, uiTemplate.getConfiguration()).build())
.addTags(TagAndValue.newLongTag(UITemplate.UPDATE_TIME, uiTemplate.getUpdateTime()).build())
.build();
} }
} }

View File

@ -0,0 +1,32 @@
/*
* 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 java.util.List;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.IndexRule;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Measure;
@RequiredArgsConstructor
@Getter
public class MeasureModel {
private final Measure measure;
private final List<IndexRule> indexRules;
}

View File

@ -48,20 +48,30 @@ import lombok.Setter;
import lombok.Singular; import lombok.Singular;
import lombok.ToString; import lombok.ToString;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery; 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.Metadata;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.IntervalRule;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.Catalog;
import org.apache.skywalking.banyandb.common.v1.BanyandbCommon.ResourceOpts;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.IndexRule;
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.TagFamilySpec;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.TagSpec;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.TagType;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.TopNAggregation;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.FieldSpec;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.FieldType;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.CompressionMethod;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.EncodingMethod;
import org.apache.skywalking.banyandb.model.v1.BanyandbModel;
import org.apache.skywalking.banyandb.v1.client.BanyanDBClient; import org.apache.skywalking.banyandb.v1.client.BanyanDBClient;
import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException; import org.apache.skywalking.banyandb.v1.client.grpc.exception.BanyanDBException;
import org.apache.skywalking.banyandb.v1.client.metadata.Catalog;
import org.apache.skywalking.banyandb.v1.client.metadata.Duration; import org.apache.skywalking.banyandb.v1.client.metadata.Duration;
import org.apache.skywalking.banyandb.v1.client.metadata.Group; import org.apache.skywalking.banyandb.v1.client.metadata.MetadataCache;
import org.apache.skywalking.banyandb.v1.client.metadata.IndexRule;
import org.apache.skywalking.banyandb.v1.client.metadata.IntervalRule;
import org.apache.skywalking.banyandb.v1.client.metadata.Measure;
import org.apache.skywalking.banyandb.v1.client.metadata.NamedSchema;
import org.apache.skywalking.banyandb.v1.client.metadata.ResourceExist; import org.apache.skywalking.banyandb.v1.client.metadata.ResourceExist;
import org.apache.skywalking.banyandb.v1.client.metadata.Stream;
import org.apache.skywalking.banyandb.v1.client.metadata.TagFamilySpec;
import org.apache.skywalking.banyandb.v1.client.metadata.TopNAggregation;
import org.apache.skywalking.oap.server.core.analysis.DownSampling; import org.apache.skywalking.oap.server.core.analysis.DownSampling;
import org.apache.skywalking.oap.server.core.analysis.metrics.IntList; import org.apache.skywalking.oap.server.core.analysis.metrics.IntList;
import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics; import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics;
@ -86,7 +96,7 @@ public enum MetadataRegistry {
private Map<String, GroupSetting> specificGroupSettings = new HashMap<>(); private Map<String, GroupSetting> specificGroupSettings = new HashMap<>();
public Stream registerStreamModel(Model model, BanyanDBStorageConfig config, ConfigService configService) { public StreamModel registerStreamModel(Model model, BanyanDBStorageConfig config, ConfigService configService) {
final SchemaMetadata schemaMetadata = parseMetadata(model, config, configService); final SchemaMetadata schemaMetadata = parseMetadata(model, config, configService);
Schema.SchemaBuilder schemaBuilder = Schema.builder().metadata(schemaMetadata); Schema.SchemaBuilder schemaBuilder = Schema.builder().metadata(schemaMetadata);
Map<String, ModelColumn> modelColumnMap = model.getColumns().stream() Map<String, ModelColumn> modelColumnMap = model.getColumns().stream()
@ -100,12 +110,12 @@ public enum MetadataRegistry {
// this can be used to build both // this can be used to build both
// 1) a list of TagFamilySpec, // 1) a list of TagFamilySpec,
// 2) a list of IndexRule, // 2) a list of IndexRule,
List<TagMetadata> tags = parseTagMetadata(model, schemaBuilder, shardingColumns); List<TagMetadata> tags = parseTagMetadata(model, schemaBuilder, shardingColumns, schemaMetadata.group);
List<TagFamilySpec> tagFamilySpecs = schemaMetadata.extractTagFamilySpec(tags, false); List<TagFamilySpec> tagFamilySpecs = schemaMetadata.extractTagFamilySpec(tags, false);
// iterate over tagFamilySpecs to save tag names // iterate over tagFamilySpecs to save tag names
for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) { for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) {
for (final TagFamilySpec.TagSpec tagSpec : tagFamilySpec.tagSpecs()) { for (final TagSpec tagSpec : tagFamilySpec.getTagsList()) {
schemaBuilder.tag(tagSpec.getTagName()); schemaBuilder.tag(tagSpec.getName());
} }
} }
String timestampColumn4Stream = model.getBanyanDBModelExtension().getTimestampColumn(); String timestampColumn4Stream = model.getBanyanDBModelExtension().getTimestampColumn();
@ -119,15 +129,18 @@ public enum MetadataRegistry {
.filter(Objects::nonNull) .filter(Objects::nonNull)
.collect(Collectors.toList()); .collect(Collectors.toList());
final Stream.Builder builder = Stream.create(schemaMetadata.getGroup(), schemaMetadata.name()); final Stream.Builder builder = Stream.newBuilder();
builder.setEntityRelativeTags(shardingColumns); builder.setMetadata(BanyandbCommon.Metadata.newBuilder().setGroup(schemaMetadata.getGroup())
builder.addTagFamilies(tagFamilySpecs); .setName(schemaMetadata.name()));
builder.addIndexes(indexRules); builder.setEntity(BanyandbDatabase.Entity.newBuilder().addAllTagNames(shardingColumns));
builder.addAllTagFamilies(tagFamilySpecs);
//builder.addIndexes(indexRules);
registry.put(schemaMetadata.name(), schemaBuilder.build()); registry.put(schemaMetadata.name(), schemaBuilder.build());
return builder.build(); return new StreamModel(builder.build(), indexRules);
} }
public Measure registerMeasureModel(Model model, BanyanDBStorageConfig config, ConfigService configService) throws StorageException { public MeasureModel registerMeasureModel(Model model, BanyanDBStorageConfig config, ConfigService configService) throws StorageException {
final SchemaMetadata schemaMetadata = parseMetadata(model, config, configService); final SchemaMetadata schemaMetadata = parseMetadata(model, config, configService);
Schema.SchemaBuilder schemaBuilder = Schema.builder().metadata(schemaMetadata); Schema.SchemaBuilder schemaBuilder = Schema.builder().metadata(schemaMetadata);
Map<String, ModelColumn> modelColumnMap = model.getColumns().stream() Map<String, ModelColumn> modelColumnMap = model.getColumns().stream()
@ -141,12 +154,12 @@ public enum MetadataRegistry {
// this can be used to build both // this can be used to build both
// 1) a list of TagFamilySpec, // 1) a list of TagFamilySpec,
// 2) a list of IndexRule, // 2) a list of IndexRule,
MeasureMetadata tagsAndFields = parseTagAndFieldMetadata(model, schemaBuilder, shardingColumns); MeasureMetadata tagsAndFields = parseTagAndFieldMetadata(model, schemaBuilder, shardingColumns, schemaMetadata.group);
List<TagFamilySpec> tagFamilySpecs = schemaMetadata.extractTagFamilySpec(tagsAndFields.tags, model.getBanyanDBModelExtension().isStoreIDTag()); List<TagFamilySpec> tagFamilySpecs = schemaMetadata.extractTagFamilySpec(tagsAndFields.tags, model.getBanyanDBModelExtension().isStoreIDTag());
// iterate over tagFamilySpecs to save tag names // iterate over tagFamilySpecs to save tag names
for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) { for (final TagFamilySpec tagFamilySpec : tagFamilySpecs) {
for (final TagFamilySpec.TagSpec tagSpec : tagFamilySpec.tagSpecs()) { for (final TagSpec tagSpec : tagFamilySpec.getTagsList()) {
schemaBuilder.tag(tagSpec.getTagName()); schemaBuilder.tag(tagSpec.getName());
} }
} }
List<IndexRule> indexRules = tagsAndFields.tags.stream() List<IndexRule> indexRules = tagsAndFields.tags.stream()
@ -155,26 +168,29 @@ public enum MetadataRegistry {
.collect(Collectors.toList()); .collect(Collectors.toList());
if (model.getBanyanDBModelExtension().isStoreIDTag()) { if (model.getBanyanDBModelExtension().isStoreIDTag()) {
indexRules.add(IndexRule.create(BanyanDBConverter.ID, IndexRule.IndexType.INVERTED)); indexRules.add(indexRule(schemaMetadata.group, BanyanDBConverter.ID));
// indexRules.add(IndexRule.create(BanyanDBConverter.ID, IndexRule.IndexType.INVERTED));
} }
final Measure.Builder builder = Measure.create(schemaMetadata.getGroup(), schemaMetadata.name(), final Measure.Builder builder = Measure.newBuilder();
downSamplingDuration(model.getDownsampling())); builder.setMetadata(BanyandbCommon.Metadata.newBuilder().setGroup(schemaMetadata.getGroup())
builder.setEntityRelativeTags(shardingColumns); .setName(schemaMetadata.name()));
builder.addTagFamilies(tagFamilySpecs); builder.setInterval(downSamplingDuration(model.getDownsampling()).format());
if (!indexRules.isEmpty()) { builder.setEntity(BanyandbDatabase.Entity.newBuilder().addAllTagNames(shardingColumns));
builder.addIndexes(indexRules); builder.addAllTagFamilies(tagFamilySpecs);
} // if (!indexRules.isEmpty()) {
// builder.addIndexes(indexRules);
// }
// parse and set field // parse and set field
for (Measure.FieldSpec field : tagsAndFields.fields) { for (BanyandbDatabase.FieldSpec field : tagsAndFields.fields) {
builder.addField(field); builder.addFields(field);
schemaBuilder.field(field.getName()); schemaBuilder.field(field.getName());
} }
// parse TopN // parse TopN
schemaBuilder.topNSpec(parseTopNSpec(model, schemaMetadata.name())); schemaBuilder.topNSpec(parseTopNSpec(model, schemaMetadata.name()));
registry.put(schemaMetadata.name(), schemaBuilder.build()); registry.put(schemaMetadata.name(), schemaBuilder.build());
return builder.build(); return new MeasureModel(builder.build(), indexRules);
} }
private TopNSpec parseTopNSpec(final Model model, final String measureName) private TopNSpec parseTopNSpec(final Model model, final String measureName)
@ -198,7 +214,7 @@ public enum MetadataRegistry {
.countersNumber(model.getBanyanDBModelExtension().getTopN().getCountersNumber()) .countersNumber(model.getBanyanDBModelExtension().getTopN().getCountersNumber())
.fieldName(valueColumnOpt.get().getValueCName()) .fieldName(valueColumnOpt.get().getValueCName())
.groupByTagNames(model.getBanyanDBModelExtension().getTopN().getGroupByTagNames()) .groupByTagNames(model.getBanyanDBModelExtension().getTopN().getGroupByTagNames())
.sort(AbstractQuery.Sort.UNSPECIFIED) // include both TopN and BottomN .sort(BanyandbModel.Sort.SORT_UNSPECIFIED) // include both TopN and BottomN
.build(); .build();
} }
@ -245,27 +261,31 @@ public enum MetadataRegistry {
return this.registry.get(SchemaMetadata.formatName(modelName, downSampling)); return this.registry.get(SchemaMetadata.formatName(modelName, downSampling));
} }
private Measure.FieldSpec parseFieldSpec(ModelColumn modelColumn) { private FieldSpec parseFieldSpec(ModelColumn modelColumn) {
String colName = modelColumn.getColumnName().getStorageName(); String colName = modelColumn.getColumnName().getStorageName();
if (String.class.equals(modelColumn.getType())) { if (String.class.equals(modelColumn.getType())) {
return Measure.FieldSpec.newIntField(colName) return FieldSpec.newBuilder().setName(colName)
.compressWithZSTD() .setFieldType(FieldType.FIELD_TYPE_STRING)
.build(); .setCompressionMethod(CompressionMethod.COMPRESSION_METHOD_ZSTD)
.build();
} else if (long.class.equals(modelColumn.getType()) || int.class.equals(modelColumn.getType())) { } else if (long.class.equals(modelColumn.getType()) || int.class.equals(modelColumn.getType())) {
return Measure.FieldSpec.newIntField(colName) return FieldSpec.newBuilder().setName(colName)
.compressWithZSTD() .setFieldType(FieldType.FIELD_TYPE_INT)
.encodeWithGorilla() .setCompressionMethod(CompressionMethod.COMPRESSION_METHOD_ZSTD)
.build(); .setEncodingMethod(EncodingMethod.ENCODING_METHOD_GORILLA)
.build();
} else if (StorageDataComplexObject.class.isAssignableFrom(modelColumn.getType()) || JsonObject.class.equals(modelColumn.getType())) { } else if (StorageDataComplexObject.class.isAssignableFrom(modelColumn.getType()) || JsonObject.class.equals(modelColumn.getType())) {
return Measure.FieldSpec.newStringField(colName) return FieldSpec.newBuilder().setName(colName)
.compressWithZSTD() .setFieldType(FieldType.FIELD_TYPE_STRING)
.build(); .setCompressionMethod(CompressionMethod.COMPRESSION_METHOD_ZSTD)
.build();
} else if (double.class.equals(modelColumn.getType())) { } else if (double.class.equals(modelColumn.getType())) {
// TODO: natively support double/float in BanyanDB // TODO: natively support double/float in BanyanDB
log.warn("Double is stored as binary"); log.warn("Double is stored as binary");
return Measure.FieldSpec.newBinaryField(colName) return FieldSpec.newBuilder().setName(colName)
.compressWithZSTD() .setFieldType(FieldType.FIELD_TYPE_DATA_BINARY)
.build(); .setCompressionMethod(CompressionMethod.COMPRESSION_METHOD_ZSTD)
.build();
} else { } else {
throw new UnsupportedOperationException(modelColumn.getType().getSimpleName() + " is not supported for field"); throw new UnsupportedOperationException(modelColumn.getType().getSimpleName() + " is not supported for field");
} }
@ -284,8 +304,11 @@ public enum MetadataRegistry {
} }
} }
IndexRule indexRule(String tagName) { IndexRule indexRule(String group, String tagName) {
return IndexRule.create(tagName, IndexRule.IndexType.INVERTED); return IndexRule.newBuilder()
.setMetadata(Metadata.newBuilder().setName(tagName).setGroup(group))
.setType(IndexRule.Type.TYPE_INVERTED).addTags(tagName).build();
//return IndexRule.create(tagName, IndexRule.IndexType.INVERTED);
} }
/** /**
@ -314,18 +337,18 @@ public enum MetadataRegistry {
* *
* @since 9.4.0 Skip {@link Record#TIME_BUCKET} * @since 9.4.0 Skip {@link Record#TIME_BUCKET}
*/ */
List<TagMetadata> parseTagMetadata(Model model, Schema.SchemaBuilder builder, List<String> shardingColumns) { List<TagMetadata> parseTagMetadata(Model model, Schema.SchemaBuilder builder, List<String> shardingColumns, String group) {
List<TagMetadata> tagMetadataList = new ArrayList<>(); List<TagMetadata> tagMetadataList = new ArrayList<>();
for (final ModelColumn col : model.getColumns()) { for (final ModelColumn col : model.getColumns()) {
final String columnStorageName = col.getColumnName().getStorageName(); final String columnStorageName = col.getColumnName().getStorageName();
if (columnStorageName.equals(Record.TIME_BUCKET)) { if (columnStorageName.equals(Record.TIME_BUCKET)) {
continue; continue;
} }
final TagFamilySpec.TagSpec tagSpec = parseTagSpec(col); final TagSpec tagSpec = parseTagSpec(col);
builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG, col.getType())); builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG, col.getType()));
String colName = col.getColumnName().getStorageName(); String colName = col.getColumnName().getStorageName();
if (!shardingColumns.contains(colName) && col.getBanyanDBExtension().shouldIndex()) { if (!shardingColumns.contains(colName) && col.getBanyanDBExtension().shouldIndex()) {
tagMetadataList.add(new TagMetadata(indexRule(tagSpec.getTagName()), tagSpec)); tagMetadataList.add(new TagMetadata(indexRule(group, tagSpec.getName()), tagSpec));
} else { } else {
tagMetadataList.add(new TagMetadata(null, tagSpec)); tagMetadataList.add(new TagMetadata(null, tagSpec));
} }
@ -339,7 +362,7 @@ public enum MetadataRegistry {
@Singular @Singular
private final List<TagMetadata> tags; private final List<TagMetadata> tags;
@Singular @Singular
private final List<Measure.FieldSpec> fields; private final List<BanyandbDatabase.FieldSpec> fields;
} }
/** /**
@ -349,7 +372,7 @@ public enum MetadataRegistry {
* *
* @since 9.4.0 Skip {@link Metrics#TIME_BUCKET} * @since 9.4.0 Skip {@link Metrics#TIME_BUCKET}
*/ */
MeasureMetadata parseTagAndFieldMetadata(Model model, Schema.SchemaBuilder builder, List<String> shardingColumns) { MeasureMetadata parseTagAndFieldMetadata(Model model, Schema.SchemaBuilder builder, List<String> shardingColumns, String group) {
// skip metric // skip metric
MeasureMetadata.MeasureMetadataBuilder result = MeasureMetadata.builder(); MeasureMetadata.MeasureMetadataBuilder result = MeasureMetadata.builder();
for (final ModelColumn col : model.getColumns()) { for (final ModelColumn col : model.getColumns()) {
@ -362,10 +385,10 @@ public enum MetadataRegistry {
result.field(parseFieldSpec(col)); result.field(parseFieldSpec(col));
continue; continue;
} }
final TagFamilySpec.TagSpec tagSpec = parseTagSpec(col); final TagSpec tagSpec = parseTagSpec(col);
builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG, col.getType())); builder.spec(columnStorageName, new ColumnSpec(ColumnType.TAG, col.getType()));
String colName = col.getColumnName().getStorageName(); String colName = col.getColumnName().getStorageName();
result.tag(new TagMetadata(!shardingColumns.contains(colName) && col.getBanyanDBExtension().shouldIndex() ? indexRule(tagSpec.getTagName()) : null, tagSpec)); result.tag(new TagMetadata(!shardingColumns.contains(colName) && col.getBanyanDBExtension().shouldIndex() ? indexRule(group, tagSpec.getName()) : null, tagSpec));
} }
return result.build(); return result.build();
@ -378,36 +401,36 @@ public enum MetadataRegistry {
* @return a typed tag spec * @return a typed tag spec
*/ */
@Nonnull @Nonnull
private TagFamilySpec.TagSpec parseTagSpec(ModelColumn modelColumn) { private TagSpec parseTagSpec(ModelColumn modelColumn) {
final Class<?> clazz = modelColumn.getType(); final Class<?> clazz = modelColumn.getType();
final String colName = modelColumn.getColumnName().getStorageName(); final String colName = modelColumn.getColumnName().getStorageName();
TagFamilySpec.TagSpec tagSpec = null; TagSpec.Builder tagSpec = TagSpec.newBuilder().setName(colName);
if (String.class.equals(clazz) || StorageDataComplexObject.class.isAssignableFrom(clazz) || JsonObject.class.equals(clazz)) { if (String.class.equals(clazz) || StorageDataComplexObject.class.isAssignableFrom(clazz) || JsonObject.class.equals(clazz)) {
tagSpec = TagFamilySpec.TagSpec.newStringTag(colName); 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)) {
tagSpec = TagFamilySpec.TagSpec.newIntTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_INT);
} else if (byte[].class.equals(clazz)) { } else if (byte[].class.equals(clazz)) {
tagSpec = TagFamilySpec.TagSpec.newBinaryTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_DATA_BINARY);
} else if (clazz.isEnum()) { } else if (clazz.isEnum()) {
tagSpec = TagFamilySpec.TagSpec.newIntTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_INT);
} else if (double.class.equals(clazz) || Double.class.equals(clazz)) { } else if (double.class.equals(clazz) || Double.class.equals(clazz)) {
// serialize double as binary // serialize double as binary
tagSpec = TagFamilySpec.TagSpec.newBinaryTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_DATA_BINARY);
} else if (IntList.class.isAssignableFrom(clazz)) { } else if (IntList.class.isAssignableFrom(clazz)) {
tagSpec = TagFamilySpec.TagSpec.newIntArrayTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_INT_ARRAY);
} else if (List.class.isAssignableFrom(clazz)) { // handle exceptions } else if (List.class.isAssignableFrom(clazz)) { // handle exceptions
ParameterizedType t = (ParameterizedType) modelColumn.getGenericType(); ParameterizedType t = (ParameterizedType) modelColumn.getGenericType();
if (String.class.equals(t.getActualTypeArguments()[0])) { if (String.class.equals(t.getActualTypeArguments()[0])) {
tagSpec = TagFamilySpec.TagSpec.newStringArrayTag(colName); tagSpec = tagSpec.setType(TagType.TAG_TYPE_STRING_ARRAY);
} }
} } else {
if (tagSpec == null) {
throw new IllegalStateException("type " + modelColumn.getType().toString() + " is not supported"); throw new IllegalStateException("type " + modelColumn.getType().toString() + " is not supported");
} }
if (modelColumn.isIndexOnly()) { if (modelColumn.isIndexOnly()) {
tagSpec.indexedOnly(); tagSpec.setIndexedOnly(true);
} }
return tagSpec; return tagSpec.build();
} }
public void initializeIntervals(String specificGroupSettingsStr) { public void initializeIntervals(String specificGroupSettingsStr) {
@ -507,7 +530,7 @@ public enum MetadataRegistry {
return modelName + "_" + downSampling.getName(); return modelName + "_" + downSampling.getName();
} }
public Optional<NamedSchema<?>> findRemoteSchema(BanyanDBClient client) throws BanyanDBException { public Optional<Object> findRemoteSchema(BanyanDBClient client) throws BanyanDBException {
try { try {
switch (kind) { switch (kind) {
case STREAM: case STREAM:
@ -526,6 +549,17 @@ public enum MetadataRegistry {
} }
} }
public MetadataCache.EntityMetadata updateRemoteSchema(BanyanDBClient client) throws BanyanDBException {
switch (kind) {
case STREAM:
return client.updateStreamMetadataCacheFromSever(this.group, this.name());
case MEASURE:
return client.updateMeasureMetadataCacheFromSever(this.group, this.name());
default:
throw new IllegalStateException("should not reach here");
}
}
private List<TagFamilySpec> extractTagFamilySpec(List<TagMetadata> tagMetadataList, boolean shouldAddID) { private List<TagFamilySpec> extractTagFamilySpec(List<TagMetadata> tagMetadataList, boolean shouldAddID) {
final String indexFamily = SchemaMetadata.this.indexFamily(); final String indexFamily = SchemaMetadata.this.indexFamily();
final String nonIndexFamily = SchemaMetadata.this.nonIndexFamily(); final String nonIndexFamily = SchemaMetadata.this.nonIndexFamily();
@ -534,10 +568,11 @@ public enum MetadataRegistry {
final List<TagFamilySpec> tagFamilySpecs = new ArrayList<>(tagMetadataMap.size()); final List<TagFamilySpec> tagFamilySpecs = new ArrayList<>(tagMetadataMap.size());
for (final Map.Entry<String, List<TagMetadata>> entry : tagMetadataMap.entrySet()) { for (final Map.Entry<String, List<TagMetadata>> entry : tagMetadataMap.entrySet()) {
final TagFamilySpec.Builder b = TagFamilySpec.create(entry.getKey()) final TagFamilySpec.Builder b = TagFamilySpec.newBuilder();
.addTagSpecs(entry.getValue().stream().map(TagMetadata::getTagSpec).collect(Collectors.toList())); b.setName(entry.getKey());
b.addAllTags(entry.getValue().stream().map(TagMetadata::getTagSpec).collect(Collectors.toList()));
if (shouldAddID && indexFamily.equals(entry.getKey())) { if (shouldAddID && indexFamily.equals(entry.getKey())) {
b.addTagSpec(TagFamilySpec.TagSpec.newStringTag(BanyanDBConverter.ID)); b.addTags(TagSpec.newBuilder().setType(TagType.TAG_TYPE_STRING).setName(BanyanDBConverter.ID));
} }
tagFamilySpecs.add(b.build()); tagFamilySpecs.add(b.build());
} }
@ -547,19 +582,31 @@ public enum MetadataRegistry {
public boolean checkResourceExistence(BanyanDBClient client) throws BanyanDBException { public boolean checkResourceExistence(BanyanDBClient client) throws BanyanDBException {
ResourceExist resourceExist; ResourceExist resourceExist;
Group.Builder gBuilder
= Group.newBuilder()
.setMetadata(Metadata.newBuilder().setName(this.group))
.setResourceOpts(ResourceOpts.newBuilder()
.setShardNum(this.shard)
.setSegmentInterval(
IntervalRule.newBuilder()
.setUnit(
IntervalRule.Unit.UNIT_DAY)
.setNum(
this.segmentIntervalDays))
.setTtl(
IntervalRule.newBuilder()
.setUnit(
IntervalRule.Unit.UNIT_DAY)
.setNum(
this.ttlDays)));
switch (kind) { switch (kind) {
case STREAM: case STREAM:
resourceExist = client.existStream(this.group, this.name()); resourceExist = client.existStream(this.group, this.name());
if (!resourceExist.hasGroup()) { if (!resourceExist.hasGroup()) {
try { try {
Group g = client.define(Group.create(this.group, Catalog.STREAM, this.shard, Group g = client.define(gBuilder.setCatalog(Catalog.CATALOG_STREAM).build());
IntervalRule.create(
IntervalRule.Unit.DAY, this.segmentIntervalDays),
IntervalRule.create(
IntervalRule.Unit.DAY, this.ttlDays)
));
if (g != null) { if (g != null) {
log.info("group {} created", g.name()); log.info("group {} created", g.getMetadata().getName());
} }
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) { if (ex.getStatus().equals(Status.Code.ALREADY_EXISTS)) {
@ -574,14 +621,15 @@ public enum MetadataRegistry {
resourceExist = client.existMeasure(this.group, this.name()); resourceExist = client.existMeasure(this.group, this.name());
try { try {
if (!resourceExist.hasGroup()) { if (!resourceExist.hasGroup()) {
Group g = client.define(Group.create(this.group, Catalog.MEASURE, this.shard, Group g = client.define(gBuilder.setCatalog(Catalog.CATALOG_MEASURE).build());
IntervalRule.create( // Group.create(this.group, Catalog.MEASURE, this.shard,
IntervalRule.Unit.DAY, this.segmentIntervalDays), // IntervalRule.create(
IntervalRule.create( // IntervalRule.Unit.DAY, this.segmentIntervalDays),
IntervalRule.Unit.DAY, this.ttlDays) // IntervalRule.create(
)); // IntervalRule.Unit.DAY, this.ttlDays)
// ));
if (g != null) { if (g != null) {
log.info("group {} created", g.name()); log.info("group {} created", g.getMetadata().getName());
} }
} }
} catch (BanyanDBException ex) { } catch (BanyanDBException ex) {
@ -637,7 +685,7 @@ public enum MetadataRegistry {
@Getter @Getter
private static class TagMetadata { private static class TagMetadata {
private final IndexRule indexRule; private final IndexRule indexRule;
private final TagFamilySpec.TagSpec tagSpec; private final TagSpec tagSpec;
boolean isIndex() { boolean isIndex() {
return this.indexRule != null; return this.indexRule != null;
@ -679,14 +727,29 @@ public enum MetadataRegistry {
} }
return; return;
} }
client.define(TopNAggregation.create(getMetadata().getGroup(), this.getTopNSpec().getName()) TopNAggregation.Builder builder
.setSourceMeasureName(getMetadata().name()) = TopNAggregation.newBuilder()
.setFieldValueSort(this.getTopNSpec().getSort()) .setMetadata(Metadata.newBuilder()
.setFieldName(this.getTopNSpec().getFieldName()) .setGroup(getMetadata().getGroup())
.setGroupByTagNames(this.getTopNSpec().getGroupByTagNames()) .setName(this.getTopNSpec().getName()))
.setCountersNumber(this.getTopNSpec().getCountersNumber())
.setLruSize(this.getTopNSpec().getLruSize()) .setSourceMeasure(Metadata.newBuilder()
.build()); .setGroup(getMetadata().getGroup())
.setName(getMetadata().name()))
.setFieldValueSort(this.getTopNSpec().getSort())
.setFieldName(this.getTopNSpec().getFieldName())
.addAllGroupByTagNames(this.getTopNSpec().getGroupByTagNames())
.setCountersNumber(this.getTopNSpec().getCountersNumber())
.setLruSize(this.getTopNSpec().getLruSize());
client.define(builder.build());
// client.define(TopNAggregation.create(getMetadata().getGroup(), this.getTopNSpec().getName())
// .setSourceMeasureName(getMetadata().name())
// .setFieldValueSort(this.getTopNSpec().getSort())
// .setFieldName(this.getTopNSpec().getFieldName())
// .setGroupByTagNames(this.getTopNSpec().getGroupByTagNames())
// .setCountersNumber(this.getTopNSpec().getCountersNumber())
// .setLruSize(this.getTopNSpec().getLruSize())
// .build());
log.info("installed TopN schema for measure {}", getMetadata().name()); log.info("installed TopN schema for measure {}", getMetadata().name());
} }
} }
@ -700,7 +763,7 @@ public enum MetadataRegistry {
@Singular @Singular
private final List<String> groupByTagNames; private final List<String> groupByTagNames;
private final String fieldName; private final String fieldName;
private final AbstractQuery.Sort sort; private final BanyandbModel.Sort sort;
private final int lruSize; private final int lruSize;
private final int countersNumber; private final int countersNumber;
} }

View File

@ -0,0 +1,32 @@
/*
* 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 java.util.List;
import lombok.Getter;
import lombok.RequiredArgsConstructor;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.Stream;
import org.apache.skywalking.banyandb.database.v1.BanyandbDatabase.IndexRule;
@RequiredArgsConstructor
@Getter
public class StreamModel {
private final Stream stream;
private final List<IndexRule> indexRules;
}

View File

@ -19,8 +19,11 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream; package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import lombok.extern.slf4j.Slf4j; 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.v1.client.TagAndValue;
import org.apache.skywalking.banyandb.v1.client.metadata.Property; 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.profiling.continuous.storage.ContinuousProfilingPolicy;
import org.apache.skywalking.oap.server.core.storage.profiling.continuous.IContinuousProfilingPolicyDAO; 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.BanyanDBStorageClient;
@ -48,10 +51,15 @@ public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO im
} }
public Property applyAll(ContinuousProfilingPolicy policy) { public Property applyAll(ContinuousProfilingPolicy policy) {
return Property.create(GROUP, ContinuousProfilingPolicy.INDEX_NAME, policy.id().build()) return Property.newBuilder()
.addTag(TagAndValue.newStringTag(ContinuousProfilingPolicy.UUID, policy.getUuid())) .setMetadata(BanyandbProperty.Metadata.newBuilder()
.addTag(TagAndValue.newStringTag(ContinuousProfilingPolicy.CONFIGURATION_JSON, policy.getConfigurationJson())) .setId(policy.id().build())
.build(); .setContainer(BanyandbCommon.Metadata.newBuilder()
.setGroup(GROUP)
.setName(ContinuousProfilingPolicy.INDEX_NAME)))
.addTags(TagAndValue.newStringTag(ContinuousProfilingPolicy.UUID, policy.getUuid()).build())
.addTags(TagAndValue.newStringTag(ContinuousProfilingPolicy.CONFIGURATION_JSON, policy.getConfigurationJson()).build())
.build();
} }
@Override @Override
@ -65,16 +73,16 @@ public class BanyanDBContinuousProfilingPolicyDAO extends AbstractBanyanDBDAO im
} }
}).filter(Objects::nonNull).map(properties -> { }).filter(Objects::nonNull).map(properties -> {
final ContinuousProfilingPolicy policy = new ContinuousProfilingPolicy(); final ContinuousProfilingPolicy policy = new ContinuousProfilingPolicy();
policy.setServiceId(properties.id()); policy.setServiceId(properties.getMetadata().getId());
for (TagAndValue<?> tag : properties.tags()) { for (BanyandbModel.Tag tag : properties.getTagsList()) {
if (tag.getTagName().equals(ContinuousProfilingPolicy.CONFIGURATION_JSON)) { TagAndValue<?> tagAndValue = TagAndValue.fromProtobuf(tag);
policy.setConfigurationJson((String) tag.getValue()); if (tagAndValue.getTagName().equals(ContinuousProfilingPolicy.CONFIGURATION_JSON)) {
} else if (tag.getTagName().equals(ContinuousProfilingPolicy.UUID)) { policy.setConfigurationJson((String) tagAndValue.getValue());
policy.setUuid((String) tag.getValue()); } else if (tagAndValue.getTagName().equals(ContinuousProfilingPolicy.UUID)) {
policy.setUuid((String) tagAndValue.getValue());
} }
} }
return policy; return policy;
}).collect(Collectors.toList()); }).collect(Collectors.toList());
} }
}
}

View File

@ -23,7 +23,7 @@ SW_AGENT_CLIENT_JS_COMMIT=af0565a67d382b683c1dbd94c379b7080db61449
SW_AGENT_CLIENT_JS_TEST_COMMIT=4f1eb1dcdbde3ec4a38534bf01dded4ab5d2f016 SW_AGENT_CLIENT_JS_TEST_COMMIT=4f1eb1dcdbde3ec4a38534bf01dded4ab5d2f016
SW_KUBERNETES_COMMIT_SHA=1335f15bf821a40a7cd71448fa805f0be265afcc SW_KUBERNETES_COMMIT_SHA=1335f15bf821a40a7cd71448fa805f0be265afcc
SW_ROVER_COMMIT=6bbd39aa701984482330d9dfb4dbaaff0527d55c SW_ROVER_COMMIT=6bbd39aa701984482330d9dfb4dbaaff0527d55c
SW_BANYANDB_COMMIT=59c396870ac2d81ec81113802d54277fe070d91b SW_BANYANDB_COMMIT=0e734c462571dcf55dbb7761211c07d8b156521e
SW_AGENT_PHP_COMMIT=3192c553002707d344bd6774cfab5bc61f67a1d3 SW_AGENT_PHP_COMMIT=3192c553002707d344bd6774cfab5bc61f67a1d3
SW_CTL_COMMIT=d5f3597733aa5217373986d776a3ee5ee8b3c468 SW_CTL_COMMIT=d5f3597733aa5217373986d776a3ee5ee8b3c468