Enhance Converter mechanism in kernel level (#8986)

* Enhance Converter mechanism in kernel level to make BanyanDB native feature more effective.
This commit is contained in:
吴晟 Wu Sheng 2022-05-04 18:55:18 +08:00 committed by GitHub
parent b2ec540ab0
commit 32b1596e5f
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
11 changed files with 48 additions and 67 deletions

View File

@ -34,6 +34,7 @@
* [Breaking Change] Replace all configurations `**_JETTY_**` to `**_REST_**`.
* Add the support eBPF profiling field into the process entity.
* E2E: fix log test miss verify LAL and metrics.
* Enhance Converter mechanism in kernel level to make BanyanDB native feature more effective.
#### UI

View File

@ -33,7 +33,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearchMatchQuery;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.ALARM;
@ -96,7 +95,7 @@ public class AlarmRecord extends Record {
record.setStartTime(((Number) converter.get(START_TIME)).longValue());
record.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());
record.setRuleName((String) converter.get(RULE_NAME));
record.setTagsRawData(converter.getWith(TAGS_RAW_DATA, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
record.setTagsRawData(converter.getBytes(TAGS_RAW_DATA));
// Don't read the TAGS as they are only for query.
return record;
}

View File

@ -31,7 +31,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.ElasticSearchMatchQuery;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
public abstract class AbstractLogRecord extends Record {
@ -126,7 +125,7 @@ public abstract class AbstractLogRecord extends Record {
record.setContentType(((Number) converter.get(CONTENT_TYPE)).intValue());
record.setContent((String) converter.get(CONTENT));
record.setTimestamp(((Number) converter.get(TIMESTAMP)).longValue());
record.setTagsRawData(converter.getWith(TAGS_RAW_DATA, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
record.setTagsRawData(converter.getBytes(TAGS_RAW_DATA));
record.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());
}

View File

@ -32,7 +32,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SuperDataset;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
@SuperDataset
@ -121,7 +120,7 @@ public class SegmentRecord extends Record {
record.setLatency(((Number) converter.get(LATENCY)).intValue());
record.setIsError(((Number) converter.get(IS_ERROR)).intValue());
record.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());
record.setDataBinary(converter.getWith(DATA_BINARY, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
record.setDataBinary(converter.getBytes(DATA_BINARY));
// Don't read the tags as they have been in the data binary already.
return record;
}

View File

@ -28,7 +28,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SuperDataset;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
@SuperDataset
@ -95,7 +94,7 @@ public class BrowserErrorLogRecord extends Record {
record.setTimestamp(((Number) converter.get(TIMESTAMP)).longValue());
record.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());
record.setErrorCategory(((Number) converter.get(ERROR_CATEGORY)).intValue());
record.setDataBinary(converter.getWith(DATA_BINARY, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
record.setDataBinary(converter.getBytes(DATA_BINARY));
return record;
}

View File

@ -28,7 +28,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.BanyanDBSharding
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EBPF_PROFILING_DATA;
@ -38,7 +37,7 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.EB
*/
@Data
@Stream(name = EBPFProfilingDataRecord.INDEX_NAME, scopeId = EBPF_PROFILING_DATA,
builder = EBPFProfilingDataRecord.Builder.class, processor = RecordStreamProcessor.class)
builder = EBPFProfilingDataRecord.Builder.class, processor = RecordStreamProcessor.class)
public class EBPFProfilingDataRecord extends Record {
public static final String INDEX_NAME = "ebpf_profiling_data";
@ -66,10 +65,10 @@ public class EBPFProfilingDataRecord extends Record {
@Override
public String id() {
return Hashing.sha256().newHasher()
.putString(scheduleId, Charsets.UTF_8)
.putString(stackIdList, Charsets.UTF_8)
.putLong(uploadTime)
.hash().toString();
.putString(scheduleId, Charsets.UTF_8)
.putString(stackIdList, Charsets.UTF_8)
.putLong(uploadTime)
.hash().toString();
}
public static class Builder implements StorageBuilder<EBPFProfilingDataRecord> {
@ -80,7 +79,7 @@ public class EBPFProfilingDataRecord extends Record {
dataTraffic.setScheduleId((String) converter.get(SCHEDULE_ID));
dataTraffic.setTaskId((String) converter.get(TASK_ID));
dataTraffic.setStackIdList((String) converter.get(STACK_ID_LIST));
dataTraffic.setStacksBinary(converter.getWith(STACKS_BINARY, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
dataTraffic.setStacksBinary(converter.getBytes(STACKS_BINARY));
dataTraffic.setStackDumpCount(((Number) converter.get(STACK_DUMP_COUNT)).longValue());
dataTraffic.setUploadTime(((Number) converter.get(UPLOAD_TIME)).longValue());
dataTraffic.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());

View File

@ -18,7 +18,6 @@
package org.apache.skywalking.oap.server.core.profiling.trace;
import java.util.Base64;
import lombok.Getter;
import lombok.Setter;
import org.apache.skywalking.oap.server.core.Const;
@ -32,7 +31,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.QueryUnifiedInde
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
import org.apache.skywalking.oap.server.library.util.StringUtil;
import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.PROFILE_TASK_SEGMENT_SNAPSHOT;
@ -81,11 +79,7 @@ public class ProfileThreadSnapshotRecord extends Record {
snapshot.setDumpTime(((Number) converter.get(DUMP_TIME)).longValue());
snapshot.setSequence(((Number) converter.get(SEQUENCE)).intValue());
snapshot.setTimeBucket(((Number) converter.get(TIME_BUCKET)).intValue());
if (StringUtil.isEmpty((String) converter.get(STACK_BINARY))) {
snapshot.setStackBinary(new byte[] {});
} else {
snapshot.setStackBinary(Base64.getDecoder().decode((String) converter.get(STACK_BINARY)));
}
snapshot.setStackBinary(converter.getBytes(STACK_BINARY));
return snapshot;
}
@ -96,7 +90,7 @@ public class ProfileThreadSnapshotRecord extends Record {
converter.accept(DUMP_TIME, storageData.getDumpTime());
converter.accept(SEQUENCE, storageData.getSequence());
converter.accept(TIME_BUCKET, storageData.getTimeBucket());
converter.accept(STACK_BINARY, new String(Base64.getEncoder().encode(storageData.getStackBinary())));
converter.accept(STACK_BINARY, storageData.getStackBinary());
}
}
}

View File

@ -18,8 +18,6 @@
package org.apache.skywalking.oap.server.core.storage.type;
import java.util.function.Function;
/**
* A function supplier to convert raw data from database to object defined in OAP
*/
@ -27,11 +25,10 @@ public interface Convert2Entity {
Object get(String fieldName);
/**
* Use the given type decoder to decode value of given field name.
* Get byte[] value of the given field.
*
* @param fieldName to read value
* @param typeDecoder to decode the value
* @return decoded value
* @param fieldName to read value
* @return byte[]
*/
<T, R> R getWith(String fieldName, Function<T, R> typeDecoder);
byte[] getBytes(String fieldName);
}

View File

@ -22,12 +22,17 @@ import java.util.Base64;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.function.Function;
import lombok.RequiredArgsConstructor;
import org.apache.skywalking.oap.server.core.Const;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.library.util.StringUtil;
/**
* HashMapConverter represents a HashMap based converter to hold the value.
*
* This converter includes following converting rules.
* 1. byte[] is converted to/from String through BASE64 encoder/decoder.
*/
public class HashMapConverter {
/**
* Stateful Hashmap based converter, build object from a HashMap type source.
@ -42,27 +47,12 @@ public class HashMapConverter {
}
@Override
public <T, R> R getWith(final String fieldName, final Function<T, R> typeDecoder) {
final T value = (T) source.get(fieldName);
return typeDecoder.apply(value);
}
/**
* Default Base64Decoder supplier
*/
public static class Base64Decoder implements Function<String, byte[]> {
public static final Base64Decoder INSTANCE = new Base64Decoder();
private Base64Decoder() {
}
@Override
public byte[] apply(final String encodedStr) {
if (StringUtil.isEmpty(encodedStr)) {
return new byte[] {};
}
return Base64.getDecoder().decode(encodedStr);
public byte[] getBytes(final String fieldName) {
final String value = (String) source.get(fieldName);
if (StringUtil.isEmpty(value)) {
return new byte[] {};
}
return Base64.getDecoder().decode(value);
}
}

View File

@ -22,7 +22,6 @@ import com.google.common.base.Splitter;
import java.util.ArrayList;
import java.util.Base64;
import java.util.List;
import java.util.function.Function;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.read.common.Field;
import org.apache.iotdb.tsfile.read.common.RowRecord;
@ -32,6 +31,7 @@ import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.StorageDataComplexObject;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.library.util.StringUtil;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBClient;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBIndexes;
import org.apache.skywalking.oap.server.storage.plugin.iotdb.IoTDBTableMetaInfo;
@ -82,21 +82,24 @@ public class IoTDBDataConverter {
if (IoTDBClient.TIMESTAMP.equals(fieldName) || fieldName.contains(".")) {
String columnName = IoTDBUtils.addQuotationMark(fieldName);
return columnNames.contains(columnName)
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(columnName) - 1))
: null;
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(columnName) - 1))
: null;
} else {
return columnNames.contains(fieldName)
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1))
: null;
? IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1))
: null;
}
}
}
@Override
public <T, R> R getWith(final String fieldName, final Function<T, R> typeDecoder) {
public byte[] getBytes(final String fieldName) {
if (columnNames.contains(fieldName)) {
final T value = (T) IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1));
return typeDecoder.apply(value);
final String value = (String) IoTDBUtils.getFieldValue(fields.get(columnNames.indexOf(fieldName) - 1));
if (StringUtil.isEmpty(value)) {
return new byte[] {};
}
return Base64.getDecoder().decode(value);
} else {
return null;
}
@ -121,8 +124,10 @@ public class IoTDBDataConverter {
List<String> indexValues = request.getIndexValues();
// To avoid indexValue be "null" when inserting, replace null to empty string
if (indexes.contains(fieldName)) {
indexValues.set(indexes.indexOf(fieldName),
fieldValue == null ? Const.EMPTY_STRING : fieldValue.toString());
indexValues.set(
indexes.indexOf(fieldName),
fieldValue == null ? Const.EMPTY_STRING : fieldValue.toString()
);
}
} else {
// time_bucket has changed to timestamp before calling this method,
@ -172,12 +177,12 @@ public class IoTDBDataConverter {
public Object get(String fieldName) {
if (IoTDBIndexes.isIndex(fieldName)) {
return request.getIndexes().contains(fieldName)
? request.getIndexValues().get(request.getIndexes().indexOf(fieldName))
: null;
? request.getIndexValues().get(request.getIndexes().indexOf(fieldName))
: null;
} else {
return request.getMeasurements().contains(fieldName)
? request.getMeasurementValues().get(request.getMeasurements().indexOf(fieldName))
: null;
? request.getMeasurementValues().get(request.getMeasurements().indexOf(fieldName))
: null;
}
}

View File

@ -30,7 +30,6 @@ import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.SuperDataset;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Entity;
import org.apache.skywalking.oap.server.core.storage.type.Convert2Storage;
import org.apache.skywalking.oap.server.core.storage.type.HashMapConverter;
import org.apache.skywalking.oap.server.core.storage.type.StorageBuilder;
@SuperDataset
@ -127,7 +126,7 @@ public class ZipkinSpanRecord extends Record {
record.setLatency(((Number) converter.get(LATENCY)).intValue());
record.setIsError(((Number) converter.get(IS_ERROR)).intValue());
record.setTimeBucket(((Number) converter.get(TIME_BUCKET)).longValue());
record.setDataBinary(converter.getWith(DATA_BINARY, HashMapConverter.ToEntity.Base64Decoder.INSTANCE));
record.setDataBinary(converter.getBytes(DATA_BINARY));
record.setEncode(((Number) converter.get(ENCODE)).intValue());
// Don't read the tags as they have been in the data binary already.
return record;