diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index f3d0704132..84c4639321 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -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 diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java index 8ca134d3fe..8de3fe88c4 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmRecord.java @@ -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; } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/log/AbstractLogRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/log/AbstractLogRecord.java index 0ea0ee01ee..6ac2678c07 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/log/AbstractLogRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/log/AbstractLogRecord.java @@ -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()); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java index d25091779c..24d41dc436 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java @@ -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; } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/browser/manual/errorlog/BrowserErrorLogRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/browser/manual/errorlog/BrowserErrorLogRecord.java index df12484b0b..88287f8653 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/browser/manual/errorlog/BrowserErrorLogRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/browser/manual/errorlog/BrowserErrorLogRecord.java @@ -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; } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingDataRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingDataRecord.java index 98b1519cef..5ff13638ff 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingDataRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/ebpf/storage/EBPFProfilingDataRecord.java @@ -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 { @@ -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()); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/trace/ProfileThreadSnapshotRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/trace/ProfileThreadSnapshotRecord.java index a66633a797..e69db8ffbd 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/trace/ProfileThreadSnapshotRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/profiling/trace/ProfileThreadSnapshotRecord.java @@ -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()); } } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/Convert2Entity.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/Convert2Entity.java index 4444d6e41c..e357e94352 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/Convert2Entity.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/Convert2Entity.java @@ -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[] */ - R getWith(String fieldName, Function typeDecoder); + byte[] getBytes(String fieldName); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/HashMapConverter.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/HashMapConverter.java index 9fc5c1ecdf..fff49fe1fa 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/HashMapConverter.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/type/HashMapConverter.java @@ -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 R getWith(final String fieldName, final Function typeDecoder) { - final T value = (T) source.get(fieldName); - return typeDecoder.apply(value); - } - - /** - * Default Base64Decoder supplier - */ - public static class Base64Decoder implements Function { - 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); } } diff --git a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBDataConverter.java b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBDataConverter.java index 60089cd793..036d377866 100644 --- a/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBDataConverter.java +++ b/oap-server/server-storage-plugin/storage-iotdb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/iotdb/utils/IoTDBDataConverter.java @@ -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 R getWith(final String fieldName, final Function 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 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; } } diff --git a/oap-server/server-storage-plugin/storage-zipkin-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/zipkin/ZipkinSpanRecord.java b/oap-server/server-storage-plugin/storage-zipkin-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/zipkin/ZipkinSpanRecord.java index 16376804f8..199a232ed1 100644 --- a/oap-server/server-storage-plugin/storage-zipkin-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/zipkin/ZipkinSpanRecord.java +++ b/oap-server/server-storage-plugin/storage-zipkin-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/zipkin/ZipkinSpanRecord.java @@ -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;