From 6a75de5ece9cf8f20d4fc2b06bb31306ff991008 Mon Sep 17 00:00:00 2001 From: Wan Kai Date: Thu, 24 Nov 2022 21:20:36 +0800 Subject: [PATCH] Add `@BanyanDB.TimestampColumn` to identify `which column in Record` is providing the timestamp(milliseconds) for BanyanDB (#10019) --- docs/en/changes/changes.md | 4 ++ .../DatabaseSlowStatementBuilder.java | 4 ++ .../parser/listener/SampledTraceBuilder.java | 5 ++- .../vservice/VirtualCacheProcessor.java | 1 + .../vservice/VirtualDatabaseProcessor.java | 1 + .../dsl/spec/extractor/ExtractorSpec.java | 1 + .../oap/server/core/alarm/AlarmRecord.java | 1 + .../cache/CacheSlowAccessDispatcher.java | 2 + .../manual/cache/TopNCacheReadCommand.java | 4 ++ .../manual/cache/TopNCacheWriteCommand.java | 4 ++ .../database/DatabaseStatementDispatcher.java | 1 + .../database/TopNDatabaseStatement.java | 4 ++ .../core/analysis/manual/log/LogRecord.java | 2 + .../manual/segment/SegmentRecord.java | 1 + .../spanattach/SpanAttachedEventRecord.java | 8 ++++ .../manual/trace/SampledSlowTraceRecord.java | 10 ++++- .../trace/SampledStatus4xxTraceRecord.java | 8 ++++ .../trace/SampledStatus5xxTraceRecord.java | 8 ++++ .../oap/server/core/analysis/topn/TopN.java | 5 +++ .../errorlog/BrowserErrorLogRecord.java | 1 + .../ebpf/storage/EBPFProfilingDataRecord.java | 3 +- .../ebpf/storage/EBPFProfilingTaskRecord.java | 1 + .../profiling/trace/ProfileTaskLogRecord.java | 8 ++++ .../profiling/trace/ProfileTaskRecord.java | 1 + .../trace/ProfileThreadSnapshotRecord.java | 1 + .../server/core/source/CacheSlowAccess.java | 3 ++ .../core/source/DatabaseSlowStatement.java | 3 ++ .../core/storage/annotation/BanyanDB.java | 12 ++++++ .../storage/model/BanyanDBModelExtension.java | 40 +++++++++++++++++++ .../oap/server/core/storage/model/Model.java | 5 ++- .../core/storage/model/StorageModels.java | 10 ++++- .../server/core/zipkin/ZipkinSpanRecord.java | 3 +- .../handler/ProfileTaskServiceHandler.java | 5 ++- ...SpanAttachedEventReportServiceHandler.java | 9 +++-- .../plugin/banyandb/BanyanDBConverter.java | 3 ++ .../banyandb/BanyanDBNoneStreamDAO.java | 10 ++--- .../banyandb/BanyanDBZipkinQueryDAO.java | 36 ++++++++--------- .../plugin/banyandb/MetadataRegistry.java | 10 +++++ .../banyandb/stream/BanyanDBRecordDAO.java | 10 ++--- .../base/TimeSeriesUtilsTest.java | 10 +++-- 40 files changed, 214 insertions(+), 44 deletions(-) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/BanyanDBModelExtension.java diff --git a/docs/en/changes/changes.md b/docs/en/changes/changes.md index e865f06135..412a766e73 100644 --- a/docs/en/changes/changes.md +++ b/docs/en/changes/changes.md @@ -119,6 +119,10 @@ * Support dynamic config the sampling strategy in network profiling. * Zipkin module support BanyanDB storage. * Zipkin traces query API, sort the result set by start time by default. +* [**Breaking Change**] Add `@BanyanDB.TimestampColumn` to identify `which column in Record` is providing the timestamp(milliseconds) for BanyanDB, + since BanyanDB stream requires a timestamp in milliseconds. + For SQL-Database: add new column `timestamp` for tables `profile_task_log/top_n_database_statement`, + requires altering this column or removing these tables before OAP starts, if bump up from previous releases. #### UI diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/DatabaseSlowStatementBuilder.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/DatabaseSlowStatementBuilder.java index 6a9e1e9d0a..47b59ab721 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/DatabaseSlowStatementBuilder.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/DatabaseSlowStatementBuilder.java @@ -51,6 +51,9 @@ public class DatabaseSlowStatementBuilder { @Getter @Setter private long timeBucket; + @Getter + @Setter + private long timestamp; public void prepare() { this.serviceName = namingControl.formatServiceName(serviceName); @@ -64,6 +67,7 @@ public class DatabaseSlowStatementBuilder { dbSlowStat.setStatement(statement); dbSlowStat.setLatency(latency); dbSlowStat.setTimeBucket(timeBucket); + dbSlowStat.setTimestamp(timestamp); return dbSlowStat; } } diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/SampledTraceBuilder.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/SampledTraceBuilder.java index 799155bba6..8a64d84eb9 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/SampledTraceBuilder.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/SampledTraceBuilder.java @@ -107,6 +107,7 @@ public class SampledTraceBuilder { slowTraceRecord.setUri(uri); slowTraceRecord.setLatency(latency); slowTraceRecord.setTimeBucket(TimeBucket.getTimeBucket(timestamp, DownSampling.Second)); + slowTraceRecord.setTimestamp(timestamp); return slowTraceRecord; case STATUS_4XX: final SampledStatus4xxTraceRecord status4xxTraceRecord = new SampledStatus4xxTraceRecord(); @@ -118,6 +119,7 @@ public class SampledTraceBuilder { status4xxTraceRecord.setUri(uri); status4xxTraceRecord.setLatency(latency); status4xxTraceRecord.setTimeBucket(TimeBucket.getTimeBucket(timestamp, DownSampling.Second)); + status4xxTraceRecord.setTimestamp(timestamp); return status4xxTraceRecord; case STATUS_5XX: final SampledStatus5xxTraceRecord status5xxTraceRecord = new SampledStatus5xxTraceRecord(); @@ -129,6 +131,7 @@ public class SampledTraceBuilder { status5xxTraceRecord.setUri(uri); status5xxTraceRecord.setLatency(latency); status5xxTraceRecord.setTimeBucket(TimeBucket.getTimeBucket(timestamp, DownSampling.Second)); + status5xxTraceRecord.setTimestamp(timestamp); return status5xxTraceRecord; default: throw new IllegalArgumentException("unknown reason: " + this.reason); @@ -156,4 +159,4 @@ public class SampledTraceBuilder { STATUS_4XX, STATUS_5XX } -} \ No newline at end of file +} diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualCacheProcessor.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualCacheProcessor.java index de67db00c2..ccf387a8fd 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualCacheProcessor.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualCacheProcessor.java @@ -89,6 +89,7 @@ public class VirtualCacheProcessor implements VirtualServiceProcessor { slowAccess.setCommand(tags.get(SpanTags.CACHE_CMD)); slowAccess.setKey(tags.get(SpanTags.CACHE_KEY)); slowAccess.setTimeBucket(TimeBucket.getRecordTimeBucket(span.getStartTime())); + slowAccess.setTimestamp(span.getStartTime()); slowAccess.setOperation(op); sourceList.add(slowAccess); } diff --git a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualDatabaseProcessor.java b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualDatabaseProcessor.java index 103303c02e..9e8c7df2cd 100644 --- a/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualDatabaseProcessor.java +++ b/oap-server/analyzer/agent-analyzer/src/main/java/org/apache/skywalking/oap/server/analyzer/provider/trace/parser/listener/vservice/VirtualDatabaseProcessor.java @@ -72,6 +72,7 @@ public class VirtualDatabaseProcessor implements VirtualServiceProcessor { dbSlowStat.setStatement(statement); dbSlowStat.setLatency(latency); dbSlowStat.setTimeBucket(TimeBucket.getRecordTimeBucket(span.getStartTime())); + dbSlowStat.setTimestamp(span.getStartTime()); recordList.add(dbSlowStat); }); } diff --git a/oap-server/analyzer/log-analyzer/src/main/java/org/apache/skywalking/oap/log/analyzer/dsl/spec/extractor/ExtractorSpec.java b/oap-server/analyzer/log-analyzer/src/main/java/org/apache/skywalking/oap/log/analyzer/dsl/spec/extractor/ExtractorSpec.java index 96c0d3fefc..d535492ef2 100644 --- a/oap-server/analyzer/log-analyzer/src/main/java/org/apache/skywalking/oap/log/analyzer/dsl/spec/extractor/ExtractorSpec.java +++ b/oap-server/analyzer/log-analyzer/src/main/java/org/apache/skywalking/oap/log/analyzer/dsl/spec/extractor/ExtractorSpec.java @@ -289,6 +289,7 @@ public class ExtractorSpec extends AbstractSpec { long timeBucketForDB = TimeBucket.getTimeBucket(log.getTimestamp(), DownSampling.Second); builder.setTimeBucket(timeBucketForDB); + builder.setTimestamp(log.getTimestamp()); String entityId = serviceMeta.getEntityId(); builder.prepare(); 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 84fe19d8cd..2a666b15ac 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 @@ -43,6 +43,7 @@ import static org.apache.skywalking.oap.server.core.source.DefaultScopeDefine.AL @ScopeDeclaration(id = ALARM, name = "Alarm") @Stream(name = AlarmRecord.INDEX_NAME, scopeId = DefaultScopeDefine.ALARM, builder = AlarmRecord.Builder.class, processor = RecordStreamProcessor.class) @SQLDatabase.ExtraColumn4AdditionalEntity(additionalTable = AlarmRecord.ADDITIONAL_TAG_TABLE, parentColumn = TIME_BUCKET) +@BanyanDB.TimestampColumn(AlarmRecord.START_TIME) public class AlarmRecord extends Record { public static final String INDEX_NAME = "alarm_record"; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/cache/CacheSlowAccessDispatcher.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/cache/CacheSlowAccessDispatcher.java index f72b6da884..9db95006a9 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/cache/CacheSlowAccessDispatcher.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/cache/CacheSlowAccessDispatcher.java @@ -36,6 +36,7 @@ public class CacheSlowAccessDispatcher implements SourceDispatcher streamClass; private final boolean timeRelativeID; private final SQLDatabaseModelExtension sqlDBModelExtension; + private final BanyanDBModelExtension banyanDBModelExtension; public Model(final String name, final List columns, @@ -48,7 +49,8 @@ public class Model { final boolean superDataset, final Class streamClass, boolean timeRelativeID, - final SQLDatabaseModelExtension sqlDBModelExtension) { + final SQLDatabaseModelExtension sqlDBModelExtension, + final BanyanDBModelExtension banyanDBModelExtension) { this.name = name; this.columns = columns; this.scopeId = scopeId; @@ -59,5 +61,6 @@ public class Model { this.streamClass = streamClass; this.timeRelativeID = timeRelativeID; this.sqlDBModelExtension = sqlDBModelExtension; + this.banyanDBModelExtension = banyanDBModelExtension; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java index 33969f60f2..f8fc461093 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/StorageModels.java @@ -59,6 +59,7 @@ public class StorageModels implements IModelManager, ModelCreator, ModelManipula List modelColumns = new ArrayList<>(); ShardingKeyChecker checker = new ShardingKeyChecker(); SQLDatabaseModelExtension sqlDBModelExtension = new SQLDatabaseModelExtension(); + BanyanDBModelExtension banyanDBModelExtension = new BanyanDBModelExtension(); retrieval(aClass, storage.getModelName(), modelColumns, scopeId, checker, sqlDBModelExtension, record); // Add extra column for additional entities if (aClass.isAnnotationPresent(SQLDatabase.ExtraColumn4AdditionalEntity.class) @@ -86,6 +87,12 @@ public class StorageModels implements IModelManager, ModelCreator, ModelManipula } }); } + //Add timestampColumn for BanyanDB + if (aClass.isAnnotationPresent(BanyanDB.TimestampColumn.class)) { + String timestampColumn = aClass.getAnnotation(BanyanDB.TimestampColumn.class).value(); + banyanDBModelExtension.setTimestampColumn(timestampColumn); + } + checker.check(storage.getModelName()); Model model = new Model( @@ -97,7 +104,8 @@ public class StorageModels implements IModelManager, ModelCreator, ModelManipula isSuperDatasetModel(aClass), aClass, storage.isTimeRelativeID(), - sqlDBModelExtension + sqlDBModelExtension, + banyanDBModelExtension ); this.followColumnNameRules(model); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java index 3e44aa1f65..2676942032 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/zipkin/ZipkinSpanRecord.java @@ -50,6 +50,7 @@ import static org.apache.skywalking.oap.server.core.analysis.record.Record.TIME_ @Stream(name = ZipkinSpanRecord.INDEX_NAME, scopeId = DefaultScopeDefine.ZIPKIN_SPAN, builder = ZipkinSpanRecord.Builder.class, processor = RecordStreamProcessor.class) @SQLDatabase.ExtraColumn4AdditionalEntity(additionalTable = ZipkinSpanRecord.ADDITIONAL_QUERY_TABLE, parentColumn = TIME_BUCKET) @SQLDatabase.Sharding(shardingAlgorithm = ShardingAlgorithm.TIME_SEC_RANGE_SHARDING_ALGORITHM, dataSourceShardingColumn = TRACE_ID, tableShardingColumn = TIME_BUCKET) +@BanyanDB.TimestampColumn(ZipkinSpanRecord.TIMESTAMP_MILLIS) public class ZipkinSpanRecord extends Record { private static final Gson GSON = new Gson(); public static final int QUERY_LENGTH = 256; @@ -167,7 +168,7 @@ public class ZipkinSpanRecord extends Record { @Override public String id() { - return spanId + Const.LINE + kind; + return traceId + Const.LINE + spanId; } public static class Builder implements StorageBuilder { diff --git a/oap-server/server-receiver-plugin/skywalking-profile-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/profile/provider/handler/ProfileTaskServiceHandler.java b/oap-server/server-receiver-plugin/skywalking-profile-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/profile/provider/handler/ProfileTaskServiceHandler.java index 44210a0fa9..1de09417ff 100644 --- a/oap-server/server-receiver-plugin/skywalking-profile-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/profile/provider/handler/ProfileTaskServiceHandler.java +++ b/oap-server/server-receiver-plugin/skywalking-profile-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/profile/provider/handler/ProfileTaskServiceHandler.java @@ -153,9 +153,10 @@ public class ProfileTaskServiceHandler extends ProfileTaskGrpc.ProfileTaskImplBa logRecord.setOperationType(operationType.getCode()); logRecord.setOperationTime(System.currentTimeMillis()); // same with task time bucket, ensure record will ttl same with profile task + long timestamp = task.getStartTime() + TimeUnit.MINUTES.toMillis(task.getDuration()); logRecord.setTimeBucket( - TimeBucket.getRecordTimeBucket(task.getStartTime() + TimeUnit.MINUTES.toMillis(task.getDuration()))); - + TimeBucket.getRecordTimeBucket(timestamp)); + logRecord.setTimestamp(timestamp); RecordStreamProcessor.getInstance().in(logRecord); } diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/SpanAttachedEventReportServiceHandler.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/SpanAttachedEventReportServiceHandler.java index 513061dca7..d7660f7d24 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/SpanAttachedEventReportServiceHandler.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/handler/v8/grpc/SpanAttachedEventReportServiceHandler.java @@ -57,9 +57,10 @@ public class SpanAttachedEventReportServiceHandler extends SpanAttachedEventRepo record.setTraceSegmentId(event.getTraceContext().getTraceSegmentId()); record.setTraceSpanId(event.getTraceContext().getSpanId()); record.setDataBinary(event.toByteArray()); - record.setTimeBucket(TimeBucket.getMinuteTimeBucket(TimeUnit.SECONDS.toMillis(record.getStartTimeSecond()) - + TimeUnit.NANOSECONDS.toMillis(record.getStartTimeNanos()))); - + long timestamp = TimeUnit.SECONDS.toMillis(record.getStartTimeSecond()) + + TimeUnit.NANOSECONDS.toMillis(record.getStartTimeNanos()); + record.setTimeBucket(TimeBucket.getMinuteTimeBucket(timestamp)); + record.setTimestamp(timestamp); RecordStreamProcessor.getInstance().in(record); } @@ -82,4 +83,4 @@ public class SpanAttachedEventReportServiceHandler extends SpanAttachedEventRepo } }; } -} \ No newline at end of file +} diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java index e40c6b7015..acab5dcb7e 100644 --- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java +++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBConverter.java @@ -70,6 +70,9 @@ public class BanyanDBConverter { @Override public void accept(String fieldName, Object fieldValue) { + if (fieldName.equals(this.schema.getTimestampColumn4Stream())) { + streamWrite.setTimestamp((long) fieldValue); + } MetadataRegistry.ColumnSpec columnSpec = this.schema.getSpec(fieldName); if (columnSpec == null) { throw new IllegalArgumentException("fail to find field[" + fieldName + "]"); diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBNoneStreamDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBNoneStreamDAO.java index 12b03a94eb..46edb436ab 100644 --- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBNoneStreamDAO.java +++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBNoneStreamDAO.java @@ -20,7 +20,6 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb; import lombok.extern.slf4j.Slf4j; import org.apache.skywalking.banyandb.v1.client.StreamWrite; -import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.config.NoneStream; import org.apache.skywalking.oap.server.core.storage.AbstractDAO; import org.apache.skywalking.oap.server.core.storage.INoneStreamDAO; @@ -45,10 +44,11 @@ public class BanyanDBNoneStreamDAO extends AbstractDAO im if (schema == null) { throw new IOException(model.getName() + " is not registered"); } - StreamWrite streamWrite = new StreamWrite(schema.getMetadata().getGroup(), // group name - schema.getMetadata().name(), // stream-name - noneStream.id(), // identity - TimeBucket.getTimestamp(noneStream.getTimeBucket(), model.getDownsampling())); // timestamp + StreamWrite streamWrite = new StreamWrite( + schema.getMetadata().getGroup(), // group name + schema.getMetadata().name(), // stream-name + noneStream.id() // identity + ); // set timestamp inside `BanyanDBConverter.StreamToStorage` Convert2Storage convert2Storage = new BanyanDBConverter.StreamToStorage(schema, streamWrite); storageBuilder.entity2Storage(noneStream, convert2Storage); getClient().write(streamWrite); diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBZipkinQueryDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBZipkinQueryDAO.java index 96560abb28..3efb6f093b 100644 --- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBZipkinQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/BanyanDBZipkinQueryDAO.java @@ -24,7 +24,6 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; import java.util.LinkedHashMap; -import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -178,35 +177,33 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk public List> getTraces(final QueryRequest request, Duration duration) throws IOException { final int tracesLimit = request.limit(); int scrollLimit = 1000; - int scrollFrom = 0; + long scrollEndTime = duration.getEndTimestamp(); Set traceIds = new HashSet<>(); while (traceIds.size() < tracesLimit) { - Set resp = getTraceIds(request, duration, scrollFrom, scrollLimit); - if (resp.size() == 0) { - break; - } - for (String traceId : resp) { - traceIds.add(traceId); + List spans = getSpans(request, duration, scrollEndTime, scrollLimit); + for (ZipkinSpanRecord span : spans) { + traceIds.add(span.getTraceId()); if (traceIds.size() >= tracesLimit) { break; } } - scrollFrom = scrollFrom + scrollLimit; + if (spans.size() < scrollLimit) { + break; + } + scrollEndTime = spans.get(spans.size() - 1).getTimestampMillis(); } return getTraces(traceIds); } - private Set getTraceIds(final QueryRequest request, + private List getSpans(final QueryRequest request, Duration duration, - int from, + long scrollEndTime, int limit) throws IOException { final long startTimeMillis = duration.getStartTimestamp(); - final long endTimeMillis = duration.getEndTimestamp(); - TimestampRange tsRange = null; - if (startTimeMillis > 0 && endTimeMillis > 0) { - tsRange = new TimestampRange(startTimeMillis, endTimeMillis); + if (startTimeMillis > 0 && scrollEndTime > 0) { + tsRange = new TimestampRange(startTimeMillis, scrollEndTime); } final QueryBuilder queryBuilder = new QueryBuilder() { @@ -242,15 +239,16 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk } query.setOrderBy(new StreamQuery.OrderBy(ZipkinSpanRecord.TIMESTAMP_MILLIS, AbstractQuery.Sort.DESC)); query.setLimit(limit); - query.setOffset(from); } }; StreamQueryResponse resp = query(ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, tsRange, queryBuilder); - Set traceIds = new LinkedHashSet<>(); //needs to keep order here + List spans = new ArrayList<>(); //needs to keep order here for (final RowEntity rowEntity : resp.getElements()) { - traceIds.add(rowEntity.getTagValue(ZipkinSpanRecord.TRACE_ID)); + ZipkinSpanRecord spanRecord = new ZipkinSpanRecord.Builder().storage2Entity( + new BanyanDBConverter.StorageToStream(ZipkinSpanRecord.INDEX_NAME, rowEntity)); + spans.add(spanRecord); } - return traceIds; + return spans; } @Override diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java index 9f6f2e4cf8..f679deca66 100644 --- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java +++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/MetadataRegistry.java @@ -61,6 +61,7 @@ import java.util.Optional; import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; +import org.apache.skywalking.oap.server.library.util.StringUtil; @Slf4j public enum MetadataRegistry { @@ -87,6 +88,12 @@ public enum MetadataRegistry { schemaBuilder.tag(tagSpec.getTagName()); } } + String timestampColumn4Stream = model.getBanyanDBModelExtension().getTimestampColumn(); + if (StringUtil.isBlank(timestampColumn4Stream)) { + throw new IllegalStateException( + "Model[stream." + model.getName() + "] miss defined @BanyanDB.TimestampColumn"); + } + schemaBuilder.timestampColumn4Stream(timestampColumn4Stream); List indexRules = tags.stream() .map(TagMetadata::getIndexRule) .filter(Objects::nonNull) @@ -517,6 +524,9 @@ public enum MetadataRegistry { @Singular private final Set fields; + @Getter + private final String timestampColumn4Stream; + public ColumnSpec getSpec(String columnName) { return this.specs.get(columnName); } diff --git a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBRecordDAO.java b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBRecordDAO.java index 11b0d9abec..ed04f7e1bf 100644 --- a/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBRecordDAO.java +++ b/oap-server/server-storage-plugin/storage-banyandb-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/banyandb/stream/BanyanDBRecordDAO.java @@ -20,7 +20,6 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream; import lombok.RequiredArgsConstructor; import org.apache.skywalking.banyandb.v1.client.StreamWrite; -import org.apache.skywalking.oap.server.core.analysis.TimeBucket; import org.apache.skywalking.oap.server.core.analysis.record.Record; import org.apache.skywalking.oap.server.core.storage.IRecordDAO; import org.apache.skywalking.oap.server.core.storage.model.Model; @@ -42,10 +41,11 @@ public class BanyanDBRecordDAO implements IRecordDAO { if (schema == null) { throw new IOException(model.getName() + " is not registered"); } - StreamWrite streamWrite = new StreamWrite(schema.getMetadata().getGroup(), // group name - model.getName(), // index-name - record.id(), // identity - TimeBucket.getTimestamp(record.getTimeBucket(), model.getDownsampling())); // timestamp + StreamWrite streamWrite = new StreamWrite( + schema.getMetadata().getGroup(), // group name + model.getName(), // index-name + record.id() // identity + ); // set timestamp inside `BanyanDBConverter.StreamToStorage` Convert2Storage convert2Storage = new BanyanDBConverter.StreamToStorage(schema, streamWrite); storageBuilder.entity2Storage(record, convert2Storage); diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/test/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/TimeSeriesUtilsTest.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/test/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/TimeSeriesUtilsTest.java index a0f45b1599..8114e68494 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/test/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/TimeSeriesUtilsTest.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/test/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/TimeSeriesUtilsTest.java @@ -23,6 +23,7 @@ import org.apache.skywalking.oap.server.core.analysis.DownSampling; import org.apache.skywalking.oap.server.core.analysis.metrics.Metrics; import org.apache.skywalking.oap.server.core.analysis.record.Record; import org.apache.skywalking.oap.server.core.query.enumeration.Step; +import org.apache.skywalking.oap.server.core.storage.model.BanyanDBModelExtension; import org.apache.skywalking.oap.server.core.storage.model.Model; import org.apache.skywalking.oap.server.core.storage.model.SQLDatabaseModelExtension; import org.junit.Assert; @@ -41,13 +42,16 @@ public class TimeSeriesUtilsTest { @Before public void prepare() { superDatasetModel = new Model("superDatasetModel", Lists.newArrayList(), - 0, DownSampling.Second, true, true, Record.class, true, new SQLDatabaseModelExtension() + 0, DownSampling.Second, true, true, Record.class, true, + new SQLDatabaseModelExtension(), new BanyanDBModelExtension() ); normalRecordModel = new Model("normalRecordModel", Lists.newArrayList(), - 0, DownSampling.Second, true, false, Record.class, true, new SQLDatabaseModelExtension() + 0, DownSampling.Second, true, false, Record.class, true, + new SQLDatabaseModelExtension(), new BanyanDBModelExtension() ); normalMetricsModel = new Model("normalMetricsModel", Lists.newArrayList(), - 0, DownSampling.Minute, false, false, Metrics.class, true, new SQLDatabaseModelExtension() + 0, DownSampling.Minute, false, false, Metrics.class, true, + new SQLDatabaseModelExtension(), new BanyanDBModelExtension() ); TimeSeriesUtils.setSUPER_DATASET_DAY_STEP(1); TimeSeriesUtils.setDAY_STEP(3);