BanyanDB: Support cold stage data query for metrics/traces/logs (#13211)

This commit is contained in:
Wan Kai 2025-04-25 17:24:02 +08:00 committed by GitHub
parent c63fe21e75
commit 03b0351861
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
52 changed files with 465 additions and 359 deletions

View File

@ -48,6 +48,10 @@ extend type Query {
estimateProcessScale(serviceId: ID!, labels: [String!]!): Long!
getTimeInfo: TimeInfo
# Get the TTL info of records
getRecordsTTL: RecordsTTL
# Get the TTL info of metrics
getMetricsTTL: MetricsTTL
}
```
@ -156,6 +160,8 @@ extend type Query {
queryBasicTraces(condition: TraceQueryCondition, debug: Boolean): TraceBrief
# Read the specific trace ID with given trace ID
queryTrace(traceId: ID!, debug: Boolean): Trace
# Only for BanyanDB, can be used to query the trace in the cold stage.
queryTraceFromColdStage(traceId: ID!, duration: Duration!, debug: Boolean): Trace
# Read the list of searchable keys
queryTraceTagAutocompleteKeys(duration: Duration!):[String!]
# Search the available value options of the given key.
@ -305,6 +311,8 @@ input Duration {
start: String!
end: String!
step: Step!
# Only for BanyanDB, the flag to query from cold stage, default is false.
coldStage: Boolean
}
enum Step {

View File

@ -8,11 +8,12 @@
* BanyanDB: Support `hot/warm/cold` stages configuration.
* Fix query continues profiling policies error when the policy is already in the cache.
* Support `hot/warm/cold` stages TTL query in the status API.
* Support `hot/warm/cold` stages TTL query in the status API and graphQL API.
* PromQL Service: traffic query support `limit` and regex match.
* Fix an edge case of HashCodeSelector(Integer#MIN_VALUE causes ArrayIndexOutOfBoundsException).
* Support Flink monitoring.
* BanyanDB: Support `@ShardingKey` for Measure tags and set to TopNAggregation group tag by default.
* BanyanDB: Support cold stage data query for metrics/traces/logs.
#### UI

View File

@ -41,7 +41,8 @@ which could be accessed through HTTP GET `http://{core restHost}:{core restPort}
| expression | The MQE query expression | Yes |
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| serviceInstance | The service instance name | No |
@ -182,22 +183,23 @@ childSpans:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/queryBasicTraces?{parameters}`.
- Parameters
| Field | Description | Required |
|-------------------|---------------------------------------------------------------|--------------------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| serviceInstance | The service instance name | No |
| endpoint | The endpoint name | No |
| minTraceDuration | The minimum duration of the trace | No |
| maxTraceDuration | The maximum duration of the trace | No |
| traceState | The state of the trace, `ALL`, `SUCCESS`, `ERROR` | Yes |
| queryOrder | The order of the query result, `BY_START_TIME`, `BY_DURATION` | Yes |
| tags | The tags of the trace, `key1=value1,key2=value2` | No |
| pageNum | The page number of the query result | Yes |
| pageSize | The page size of the query result | Yes |
| Field | Description | Required |
|--------------------|---------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| serviceInstance | The service instance name | No |
| endpoint | The endpoint name | No |
| minTraceDuration | The minimum duration of the trace | No |
| maxTraceDuration | The maximum duration of the trace | No |
| traceState | The state of the trace, `ALL`, `SUCCESS`, `ERROR` | Yes |
| queryOrder | The order of the query result, `BY_START_TIME`, `BY_DURATION` | Yes |
| tags | The tags of the trace, `key1=value1,key2=value2` | No |
| pageNum | The page number of the query result | Yes |
| pageSize | The page size of the query result | Yes |
The time and step parameters are follow the [Duration](../api/query-protocol.md#duration) format.
@ -236,6 +238,19 @@ debuggingTrace:
...
```
#### Tracing SkyWalking API queryTraceFromColdStage
Only for BanyanDB, can be used to query the trace in the cold stage.
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/trace/queryTraceFromColdStage?{parameters}`.
- Parameters
| Field | Description | Required |
|-----------------|----------------------------------|-----------------|
| traceId | The ID of the trace | Yes |
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
### Tracing Zipkin Trace Query
#### Tracing Zipkin API /api/v2/traces
@ -299,12 +314,13 @@ debuggingTrace:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/topology/getGlobalTopology?{parameters}`.
- Parameters
| Field | Description | Required |
|-------------------|---------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| serviceLayer | The service layer name | No |
| Field | Description | Required |
|---------------|-------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| serviceLayer | The service layer name | No |
- Example
```shell
@ -326,13 +342,14 @@ debuggingTrace:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/topology/getServicesTopology?{parameters}`.
- Parameters
| Field | Description | Required |
|--------------|-----------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| serviceLayer | The service layer name | Yes |
| services | The services names list, separate by comma `mock_a_service, mock_b_service` | Yes |
| Field | Description | Required |
|---------------|-----------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| serviceLayer | The service layer name | Yes |
| services | The services names list, separate by comma `mock_a_service, mock_b_service` | Yes |
- Example
```shell
@ -354,15 +371,16 @@ debuggingTrace:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/topology/getServiceInstanceTopology?{parameters}`.
- Parameters
| Field | Description | Required |
|--------------------|------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| clientService | The client side service name | Yes |
| serverService | The server side service name | Yes |
| clientServiceLayer | The client side service layer name | Yes |
| serverServiceLayer | The server side service layer name | Yes |
| Field | Description | Required |
|--------------------|-------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| clientService | The client side service name | Yes |
| serverService | The server side service name | Yes |
| clientServiceLayer | The client side service layer name | Yes |
| serverServiceLayer | The server side service layer name | Yes |
- Example
```shell
@ -384,14 +402,15 @@ debuggingTrace:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/topology/getEndpointDependencies?{parameters}`.
- Parameters
| Field | Description | Required |
|--------------|-----------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| endpoint | The endpoint name | Yes |
| Field | Description | Required |
|----------------|-------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| endpoint | The endpoint name | Yes |
- Example
- Example
@ -414,14 +433,15 @@ debuggingTrace:
- URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/topology/getProcessTopology?{parameters}`.
- Parameters
| Field | Description | Required |
|---------------|-----------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| instance | The instance name | Yes |
| Field | Description | Required |
|--------------|-------------------------------------------------------------------------|----------|
| startTime | The start time of the query | Yes |
| endTime | The end time of the query | Yes |
| step | The query step | Yes |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| service | The service name | Yes |
| serviceLayer | The service layer name | Yes |
| instance | The instance name | Yes |
- Example
```shell
@ -445,24 +465,25 @@ debuggingTrace:
URL: HTTP GET `http://{core restHost}:{core restPort}/debugging/query/log/queryLogs?{parameters}`.
Parameters
| Field | Description | Required |
|----------------------------|----------------------------------------------------------------|--------------------------------|
| startTime | The start time of the query | Yes, unless traceId not empty |
| endTime | The end time of the query | Yes, unless traceId not empty |
| step | The query step | Yes, unless traceId not empty |
| service | The service name | No, require serviceLayer |
| serviceLayer | The service layer name | No |
| serviceInstance | The service instance name | No, require service |
| endpoint | The endpoint name | No, require service |
| traceId | The trace ID | No |
| segmentId | The segment ID | No, require traceId |
| spanId | The span ID | No, require traceId |
| queryOrder | The order of the query result, `ASC`, `DES` | No, default `DES` |
| tags | The tags of the trace, `key1=value1,key2=value2` | No |
| pageNum | The page number of the query result | Yes |
| pageSize | The page size of the query result | Yes |
| keywordsOfContent | The keywords of the log content, `keyword1,keyword2` | No |
| excludingKeywordsOfContent | The excluding keywords of the log content, `keyword1,keyword2` | No |
| Field | Description | Required |
|----------------------------|-------------------------------------------------------------------------|-------------------------------|
| startTime | The start time of the query | Yes, unless traceId not empty |
| endTime | The end time of the query | Yes, unless traceId not empty |
| step | The query step | Yes, unless traceId not empty |
| coldStage | Only for BanyanDB, the flag to query from cold stage, default is false. | No |
| service | The service name | No, require serviceLayer |
| serviceLayer | The service layer name | No |
| serviceInstance | The service instance name | No, require service |
| endpoint | The endpoint name | No, require service |
| traceId | The trace ID | No |
| segmentId | The segment ID | No, require traceId |
| spanId | The span ID | No, require traceId |
| queryOrder | The order of the query result, `ASC`, `DES` | No, default `DES` |
| tags | The tags of the trace, `key1=value1,key2=value2` | No |
| pageNum | The page number of the query result | Yes |
| pageSize | The page size of the query result | Yes |
| keywordsOfContent | The keywords of the log content, `keyword1,keyword2` | No |
| excludingKeywordsOfContent | The excluding keywords of the log content, `keyword1,keyword2` | No |
- Example
```shell

View File

@ -188,13 +188,13 @@ public class ProfileTaskQueryService implements Service {
public List<SegmentRecord> getTaskSegments(String taskId) throws IOException {
final List<String> profiledSegmentIdList = getProfileThreadSnapshotQueryDAO().queryProfiledSegmentIdList(taskId);
return getTraceQueryDAO().queryBySegmentIdList(profiledSegmentIdList);
return getTraceQueryDAO().queryBySegmentIdList(profiledSegmentIdList, null);
}
public List<ProfiledTraceSegments> getProfileTaskSegments(String taskId) throws IOException {
// query all profiled segments
final List<String> profiledSegmentIdList = getProfileThreadSnapshotQueryDAO().queryProfiledSegmentIdList(taskId);
final List<SegmentRecord> segmentRecords = getTraceQueryDAO().queryBySegmentIdList(profiledSegmentIdList);
final List<SegmentRecord> segmentRecords = getTraceQueryDAO().queryBySegmentIdList(profiledSegmentIdList, null);
if (CollectionUtils.isEmpty(segmentRecords)) {
return Collections.emptyList();
}
@ -215,7 +215,7 @@ public class ProfileTaskQueryService implements Service {
}
final List<SegmentRecord> traceRelatedSegments = getTraceQueryDAO().queryByTraceIdWithInstanceId(
new ArrayList<>(traceIdList),
new ArrayList<>(instanceIdList));
new ArrayList<>(instanceIdList), null);
// group by the traceId + service instanceId
final Map<String, List<SegmentRecord>> instanceTraceWithSegments = traceRelatedSegments.stream().filter(s -> {

View File

@ -28,6 +28,7 @@ import java.util.List;
import java.util.Objects;
import com.google.protobuf.InvalidProtocolBufferException;
import javax.annotation.Nullable;
import org.apache.commons.lang3.StringUtils;
import org.apache.skywalking.apm.network.common.v3.KeyIntValuePair;
import org.apache.skywalking.apm.network.common.v3.KeyStringValuePair;
@ -142,7 +143,10 @@ public class TraceQueryService implements Service {
}
}
public Trace queryTrace(final String traceId) throws IOException {
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
public Trace queryTrace(final String traceId, @Nullable final Duration duration) throws IOException {
DebuggingTraceContext traceContext = TRACE_CONTEXT.get();
DebuggingSpan span = null;
try {
@ -152,7 +156,7 @@ public class TraceQueryService implements Service {
msg.append("Condition: TraceId: ").append(traceId);
span.setMsg(msg.toString());
}
return invokeQueryTrace(traceId);
return invokeQueryTrace(traceId, duration);
} finally {
if (traceContext != null && span != null) {
traceContext.stopSpan(span);
@ -160,10 +164,10 @@ public class TraceQueryService implements Service {
}
}
private Trace invokeQueryTrace(final String traceId) throws IOException {
private Trace invokeQueryTrace(final String traceId, @Nullable final Duration duration) throws IOException {
Trace trace = new Trace();
List<SegmentRecord> segmentRecords = getTraceQueryDAO().queryByTraceIdDebuggable(traceId);
List<SegmentRecord> segmentRecords = getTraceQueryDAO().queryByTraceIdDebuggable(traceId, duration);
if (segmentRecords.isEmpty()) {
trace.getSpans().addAll(getTraceQueryDAO().doFlexibleTraceQuery(traceId));
} else {
@ -192,7 +196,7 @@ public class TraceQueryService implements Service {
if (CollectionUtils.isNotEmpty(sortedSpans)) {
final List<SpanAttachedEventRecord> spanAttachedEvents = getSpanAttachedEventQueryDAO().
querySpanAttachedEventsDebuggable(SpanAttachedEventTraceType.SKYWALKING, Arrays.asList(traceId));
querySpanAttachedEventsDebuggable(SpanAttachedEventTraceType.SKYWALKING, Arrays.asList(traceId), duration);
appendAttachedEventsToSpanDebuggable(sortedSpans, spanAttachedEvents);
}

View File

@ -34,6 +34,7 @@ public class Duration {
private String start;
private String end;
private Step step;
private boolean coldStage = false;
/**
* See {@link DurationUtils#convertToTimeBucket(Step, String)}

View File

@ -18,8 +18,10 @@
package org.apache.skywalking.oap.server.core.storage.query;
import javax.annotation.Nullable;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventTraceType;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingSpan;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingTraceContext;
import org.apache.skywalking.oap.server.library.module.Service;
@ -28,7 +30,10 @@ import java.io.IOException;
import java.util.List;
public interface ISpanAttachedEventQueryDAO extends Service {
default List<SpanAttachedEventRecord> querySpanAttachedEventsDebuggable(SpanAttachedEventTraceType type, List<String> traceIds) throws IOException {
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
default List<SpanAttachedEventRecord> querySpanAttachedEventsDebuggable(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) throws IOException {
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
DebuggingSpan span = null;
try {
@ -41,7 +46,7 @@ public interface ISpanAttachedEventQueryDAO extends Service {
.append(traceIds);
span.setMsg(builder.toString());
}
return querySpanAttachedEvents(type, traceIds);
return querySpanAttachedEvents(type, traceIds, duration);
} finally {
if (traceContext != null && span != null) {
traceContext.stopSpan(span);
@ -49,5 +54,8 @@ public interface ISpanAttachedEventQueryDAO extends Service {
}
}
List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) throws IOException;
}

View File

@ -20,6 +20,7 @@ package org.apache.skywalking.oap.server.core.storage.query;
import java.io.IOException;
import java.util.List;
import javax.annotation.Nullable;
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
import org.apache.skywalking.oap.server.core.analysis.manual.segment.SegmentRecord;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -88,7 +89,10 @@ public interface ITraceQueryDAO extends Service {
}
}
default List<SegmentRecord> queryByTraceIdDebuggable(String traceId) throws IOException {
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
default List<SegmentRecord> queryByTraceIdDebuggable(String traceId, @Nullable Duration duration) throws IOException {
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
DebuggingSpan span = null;
try {
@ -99,7 +103,7 @@ public interface ITraceQueryDAO extends Service {
.append(traceId);
span.setMsg(builder.toString());
}
return queryByTraceId(traceId);
return queryByTraceId(traceId, duration);
} finally {
if (traceContext != null && span != null) {
traceContext.stopSpan(span);
@ -120,14 +124,23 @@ public interface ITraceQueryDAO extends Service {
QueryOrder queryOrder,
final List<Tag> tags) throws IOException;
List<SegmentRecord> queryByTraceId(String traceId) throws IOException;
List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList) throws IOException;
List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<SegmentRecord> queryByTraceId(String traceId, @Nullable Duration duration) throws IOException;
/**
* This method gives more flexible for 3rd trace without segment concept, which can't search data through {@link #queryByTraceId(String)}
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList, @Nullable Duration duration) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList, @Nullable Duration duration) throws IOException;
/**
* This method gives more flexible for 3rd trace without segment concept, which can't search data through {@link #queryByTraceId(String, Duration)}
*/
List<Span> doFlexibleTraceQuery(String traceId) throws IOException;
}

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.core.storage.query;
import javax.annotation.Nullable;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingSpan;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingTraceContext;
@ -52,7 +53,10 @@ public interface IZipkinQueryDAO extends DAO {
}
}
default List<Span> getTraceDebuggable(final String traceId) throws IOException {
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
default List<Span> getTraceDebuggable(final String traceId, @Nullable final Duration duration) throws IOException {
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
DebuggingSpan span = null;
try {
@ -60,7 +64,7 @@ public interface IZipkinQueryDAO extends DAO {
span = traceContext.createSpan("Query Dao: getTrace");
span.setMsg("Condition: TraceId: " + traceId);
}
return getTrace(traceId);
return getTrace(traceId, duration);
} finally {
if (traceContext != null && span != null) {
traceContext.stopSpan(span);
@ -74,9 +78,18 @@ public interface IZipkinQueryDAO extends DAO {
List<String> getSpanNames(final String serviceName) throws IOException;
List<Span> getTrace(final String traceId) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<Span> getTrace(final String traceId, @Nullable final Duration duration) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<List<Span>> getTraces(final QueryRequest request, final Duration duration) throws IOException;
List<List<Span>> getTraces(final Set<String> traceIds) throws IOException;
/**
* @param duration nullable unless for BanyanDB query from cold stage
*/
List<List<Span>> getTraces(final Set<String> traceIds, @Nullable final Duration duration) throws IOException;
}

View File

@ -28,12 +28,15 @@ import org.apache.skywalking.oap.query.graphql.type.TimeInfo;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.IDManager;
import org.apache.skywalking.oap.server.core.query.MetadataQueryService;
import org.apache.skywalking.oap.server.core.query.TTLStatusQuery;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.query.type.Endpoint;
import org.apache.skywalking.oap.server.core.query.type.EndpointInfo;
import org.apache.skywalking.oap.server.core.query.type.Process;
import org.apache.skywalking.oap.server.core.query.type.Service;
import org.apache.skywalking.oap.server.core.query.type.ServiceInstance;
import org.apache.skywalking.oap.server.core.storage.ttl.MetricsTTL;
import org.apache.skywalking.oap.server.core.storage.ttl.RecordsTTL;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import static org.apache.skywalking.oap.query.graphql.AsyncQueryUtils.queryAsync;
@ -47,6 +50,7 @@ public class MetadataQueryV2 implements GraphQLQueryResolver {
private final ModuleManager moduleManager;
private MetadataQueryService metadataQueryService;
private TTLStatusQuery ttlStatusQuery;
public MetadataQueryV2(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
@ -61,6 +65,15 @@ public class MetadataQueryV2 implements GraphQLQueryResolver {
return metadataQueryService;
}
private TTLStatusQuery getTTLStatusQuery() {
if (ttlStatusQuery == null) {
ttlStatusQuery = moduleManager.find(CoreModule.NAME)
.provider()
.getService(TTLStatusQuery.class);
}
return ttlStatusQuery;
}
public CompletableFuture<Set<String>> listLayers() {
return queryAsync(() -> getMetadataQueryService().listLayers());
}
@ -115,4 +128,12 @@ public class MetadataQueryV2 implements GraphQLQueryResolver {
timeInfo.setTimezone(timezoneFormat.format(date));
return timeInfo;
}
public RecordsTTL getRecordsTTL() {
return getTTLStatusQuery().getTTL().getRecords();
}
public MetricsTTL getMetricsTTL() {
return getTTLStatusQuery().getTTL().getMetrics();
}
}

View File

@ -117,7 +117,28 @@ public class TraceQuery implements GraphQLQueryResolver {
DebuggingTraceContext.TRACE_CONTEXT.set(traceContext);
DebuggingSpan span = traceContext.createSpan("Query trace");
try {
Trace trace = getQueryService().queryTrace(traceId);
Trace trace = getQueryService().queryTrace(traceId, null);
if (debug) {
trace.setDebuggingTrace(traceContext.getExecTrace());
}
return trace;
} finally {
traceContext.stopSpan(span);
traceContext.stopTrace();
TRACE_CONTEXT.remove();
}
});
}
public CompletableFuture<Trace> queryTraceFromColdStage(final String traceId, Duration duration, boolean debug) {
duration.setColdStage(true);
return queryAsync(() -> {
DebuggingTraceContext traceContext = new DebuggingTraceContext(
"TraceId: " + traceId, debug, false);
DebuggingTraceContext.TRACE_CONTEXT.set(traceContext);
DebuggingSpan span = traceContext.createSpan("Query trace from cold stage");
try {
Trace trace = getQueryService().queryTrace(traceId, duration);
if (debug) {
trace.setDebuggingTrace(traceContext.getExecTrace());
}

@ -1 +1 @@
Subproject commit a9ed9eef09cf97256df9a33eab91fca1ba13096e
Subproject commit 23baed2234e4bbc18cd7ec7d47bfe7d4bc8ef363

View File

@ -120,6 +120,7 @@ public class DebuggingHTTPHandler {
@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("service") String service,
@Param("serviceLayer") String serviceLayer,
@Param("serviceInstance") Optional<String> serviceInstance,
@ -148,6 +149,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
ExpressionResult expressionResult = mqeQuery.execExpression(expression, entity, duration, true, dumpStorageRsp).join();
DebuggingTrace execTrace = expressionResult.getDebuggingTrace();
DebuggingMQERsp result = new DebuggingMQERsp(
@ -168,6 +170,7 @@ public class DebuggingHTTPHandler {
@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("minTraceDuration") Optional<Integer> minDuration,
@Param("maxTraceDuration") Optional<Integer> maxDuration,
@Param("traceState") String traceState,
@ -185,6 +188,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
Pagination pagination = new Pagination();
pagination.setPageNum(pageNum);
pagination.setPageSize(pageSize);
@ -226,6 +230,26 @@ public class DebuggingHTTPHandler {
return transToYAMLString(result);
}
/**
* Only for BanyanDB, can be used to query the trace in the cold stage.
*/
@SneakyThrows
@Get("/debugging/query/trace/queryTraceFromColdStage")
public String queryTraceFromColdStage(@Param("traceId") String traceId,
@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step) {
Duration duration = new Duration();
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
duration.setColdStage(true);
Trace trace = traceQuery.queryTraceFromColdStage(traceId, duration, true).join();
DebuggingQueryTraceRsp result = new DebuggingQueryTraceRsp(
trace.getSpans(), transformTrace(trace.getDebuggingTrace()));
return transToYAMLString(result);
}
@SneakyThrows
@Get("/debugging/query/zipkin/api/v2/traces")
public String queryZipkinTraces(@Param("serviceName") Optional<String> serviceName,
@ -293,11 +317,13 @@ public class DebuggingHTTPHandler {
public String getGlobalTopology(@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("serviceLayer") Optional<String> serviceLayer) {
Duration duration = new Duration();
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
Topology topology = topologyQuery.getGlobalTopology(duration, serviceLayer.orElse(null), true).join();
DebuggingQueryServiceTopologyRsp result = new DebuggingQueryServiceTopologyRsp(
topology.getNodes(), topology.getCalls(), transformTrace(topology.getDebuggingTrace()));
@ -309,13 +335,14 @@ public class DebuggingHTTPHandler {
public String getServicesTopology(@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("serviceLayer") String serviceLayer,
@Param("services") String services) {
Duration duration = new Duration();
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
List<String> ids = Arrays.stream(services.split(Const.COMMA))
.map(name -> IDManager.ServiceID.buildId(name, Layer.nameOf(serviceLayer).isNormal()))
.collect(Collectors.toList());
@ -330,6 +357,7 @@ public class DebuggingHTTPHandler {
public String getServiceInstanceTopology(@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("clientService") String clientService,
@Param("serverService") String serverService,
@Param("clientServiceLayer") String clientServiceLayer,
@ -338,6 +366,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
String clientServiceId = IDManager.ServiceID.buildId(clientService, Layer.nameOf(clientServiceLayer).isNormal());
String serverServiceId = IDManager.ServiceID.buildId(serverService, Layer.nameOf(serverServiceLayer).isNormal());
ServiceInstanceTopology topology = topologyQuery.getServiceInstanceTopology(clientServiceId, serverServiceId, duration, true).join();
@ -351,6 +380,7 @@ public class DebuggingHTTPHandler {
public String getEndpointDependencies(@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("service") String service,
@Param("serviceLayer") String serviceLayer,
@Param("endpoint") String endpoint) {
@ -358,6 +388,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
String endpointId = IDManager.EndpointID.buildId(
IDManager.ServiceID.buildId(service, Layer.nameOf(serviceLayer).isNormal()), endpoint);
EndpointTopology topology = topologyQuery.getEndpointDependencies(endpointId, duration, true).join();
@ -371,6 +402,7 @@ public class DebuggingHTTPHandler {
public String getProcessTopology(@Param("startTime") String startTime,
@Param("endTime") String endTime,
@Param("step") String step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("service") String service,
@Param("serviceLayer") String serviceLayer,
@Param("instance") String process) {
@ -378,6 +410,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime);
duration.setEnd(endTime);
duration.setStep(Step.valueOf(step));
coldStage.ifPresent(duration::setColdStage);
String instanceId = IDManager.ServiceInstanceID.buildId(
IDManager.ServiceID.buildId(service, Layer.nameOf(serviceLayer).isNormal()), process);
ProcessTopology topology = topologyQuery.getProcessTopology(instanceId, duration, true).join();
@ -395,6 +428,7 @@ public class DebuggingHTTPHandler {
@Param("startTime") Optional<String> startTime,
@Param("endTime") Optional<String> endTime,
@Param("step") Optional<String> step,
@Param("coldStage") Optional<Boolean> coldStage,
@Param("traceId") Optional<String> traceId,
@Param("segmentId") Optional<String> segmentId,
@Param("spanId") Optional<Integer> spanId,
@ -420,6 +454,7 @@ public class DebuggingHTTPHandler {
duration.setStart(startTime.get());
duration.setEnd(endTime.get());
duration.setStep(Step.valueOf(step.get()));
coldStage.ifPresent(duration::setColdStage);
condition.setQueryDuration(duration);
}

View File

@ -187,12 +187,12 @@ public class ZipkinQueryHandler {
if (StringUtil.isEmpty(traceId)) {
return AggregatedHttpResponse.of(BAD_REQUEST, ANY_TEXT_TYPE, "traceId is empty or null");
}
List<Span> trace = getZipkinQueryDAO().getTraceDebuggable(Span.normalizeTraceId(traceId.trim()));
List<Span> trace = getZipkinQueryDAO().getTraceDebuggable(Span.normalizeTraceId(traceId.trim()), null);
if (CollectionUtils.isEmpty(trace)) {
return AggregatedHttpResponse.of(NOT_FOUND, ANY_TEXT_TYPE, traceId + " not found");
}
appendEventsDebuggable(trace, getSpanAttachedEventQueryDAO().querySpanAttachedEventsDebuggable(
SpanAttachedEventTraceType.ZIPKIN, Arrays.asList(traceId)));
SpanAttachedEventTraceType.ZIPKIN, Arrays.asList(traceId), null));
return response(SpanBytesEncoder.JSON_V2.encodeList(trace));
} finally {
if (traceContext != null && debuggingSpan != null) {
@ -266,7 +266,7 @@ public class ZipkinQueryHandler {
}
}
List<List<Span>> traces = getZipkinQueryDAO().getTraces(normalizeTraceIds);
List<List<Span>> traces = getZipkinQueryDAO().getTraces(normalizeTraceIds, null);
appendEventsToTraces(traces);
return response(encodeTraces(traces));
}
@ -365,7 +365,7 @@ public class ZipkinQueryHandler {
}
final List<SpanAttachedEventRecord> records = getSpanAttachedEventQueryDAO().querySpanAttachedEventsDebuggable(SpanAttachedEventTraceType.ZIPKIN,
new ArrayList<>(traceIdWithSpans.keySet()));
new ArrayList<>(traceIdWithSpans.keySet()), null);
final Map<String, List<SpanAttachedEventRecord>> traceEvents = records.stream().collect(Collectors.groupingBy(SpanAttachedEventRecord::getRelatedTraceId));
for (Map.Entry<String, List<SpanAttachedEventRecord>> entry : traceEvents.entrySet()) {
appendEventsDebuggable(traceIdWithSpans.get(entry.getKey()), entry.getValue());

View File

@ -52,8 +52,8 @@ public class BanyanDBAggregationQueryDAO extends AbstractBanyanDBDAO implements
@Override
public List<SelectedRecord> sortMetrics(TopNCondition condition, String valueColumnName, Duration duration, List<KeyValue> additionalConditions) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
final String modelName = condition.getName();
final TimestampRange timestampRange = new TimestampRange(duration.getStartTimestamp(), duration.getEndTimestamp());
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(modelName, duration.getStep());
if (schema == null) {
throw new IOException("schema is not registered");
@ -71,20 +71,21 @@ public class BanyanDBAggregationQueryDAO extends AbstractBanyanDBDAO implements
if (CollectionUtils.isEmpty(additionalConditions) ||
additionalConditions.stream().map(KeyValue::getKey).collect(Collectors.toSet())
.equals(ImmutableSet.copyOf(schema.getTopNSpec().getGroupByTagNamesList()))) {
return serverSideTopN(condition, schema, spec, timestampRange, additionalConditions);
return serverSideTopN(isColdStage, condition, schema, spec, getTimestampRange(duration), additionalConditions);
}
}
return directMetricsTopN(condition, schema, valueColumnName, spec, timestampRange, additionalConditions);
return directMetricsTopN(isColdStage, condition, schema, valueColumnName, spec, getTimestampRange(duration), additionalConditions);
}
List<SelectedRecord> serverSideTopN(TopNCondition condition, MetadataRegistry.Schema schema, MetadataRegistry.ColumnSpec valueColumnSpec,
//todo: query cold stage
List<SelectedRecord> serverSideTopN(boolean isColdStage, TopNCondition condition, MetadataRegistry.Schema schema, MetadataRegistry.ColumnSpec valueColumnSpec,
TimestampRange timestampRange, List<KeyValue> additionalConditions) throws IOException {
TopNQueryResponse resp = null;
if (condition.getOrder() == Order.DES) {
resp = topNQueryDebuggable(schema, timestampRange, condition.getTopN(), AbstractQuery.Sort.DESC, additionalConditions, condition.getAttributes());
resp = topNQueryDebuggable(isColdStage, schema, timestampRange, condition.getTopN(), AbstractQuery.Sort.DESC, additionalConditions, condition.getAttributes());
} else {
resp = topNQueryDebuggable(schema, timestampRange, condition.getTopN(), AbstractQuery.Sort.ASC, additionalConditions, condition.getAttributes());
resp = topNQueryDebuggable(isColdStage, schema, timestampRange, condition.getTopN(), AbstractQuery.Sort.ASC, additionalConditions, condition.getAttributes());
}
if (resp.size() == 0) {
return Collections.emptyList();
@ -103,9 +104,9 @@ public class BanyanDBAggregationQueryDAO extends AbstractBanyanDBDAO implements
return topNList;
}
List<SelectedRecord> directMetricsTopN(TopNCondition condition, MetadataRegistry.Schema schema, String valueColumnName, MetadataRegistry.ColumnSpec valueColumnSpec,
List<SelectedRecord> directMetricsTopN(boolean isColdStage, TopNCondition condition, MetadataRegistry.Schema schema, String valueColumnName, MetadataRegistry.ColumnSpec valueColumnSpec,
TimestampRange timestampRange, List<KeyValue> additionalConditions) throws IOException {
MeasureQueryResponse resp = queryDebuggable(schema, TAGS, Collections.singleton(valueColumnName),
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema, TAGS, Collections.singleton(valueColumnName),
timestampRange, new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {

View File

@ -324,10 +324,12 @@ public class BanyanDBIndexInstaller extends ModelInstaller {
stage.getTtl()))
.setNodeSelector(stage.getNodeSelector())
.setClose(stage.isClose())
//todo: set the default query stages
);
}
}
if (CollectionUtils.isNotEmpty(metadata.getResource().getDefaultQueryStages())) {
optsBuilder.addAllDefaultStages(metadata.getResource().getDefaultQueryStages());
}
gBuilder.setResourceOpts(optsBuilder.build());
if (!RunningMode.isNoInitMode()) {
if (!groupAligned.contains(metadata.getGroup())) {

View File

@ -23,7 +23,6 @@ import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.topn.TopN;
import org.apache.skywalking.oap.server.core.query.enumeration.Order;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -47,11 +46,11 @@ public class BanyanDBRecordsQueryDAO extends AbstractBanyanDBDAO implements IRec
@Override
public List<Record> readRecords(RecordCondition condition, String valueColumnName, Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
final String modelName = condition.getName();
final TimestampRange timestampRange = new TimestampRange(duration.getStartTimestamp(), duration.getEndTimestamp());
final Set<String> tags = ImmutableSet.of(TopN.ENTITY_ID, TopN.STATEMENT, TopN.TRACE_ID, valueColumnName);
StreamQueryResponse resp = queryDebuggable(modelName, tags,
timestampRange, new QueryBuilder<StreamQuery>() {
StreamQueryResponse resp = queryDebuggable(isColdStage, modelName, tags,
getTimestampRange(duration), new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
query.and(eq(TopN.ENTITY_ID, condition.getParentEntity().buildId()));

View File

@ -27,6 +27,7 @@ import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
import javax.annotation.Nullable;
import org.apache.skywalking.banyandb.v1.client.AbstractCriteria;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
@ -87,7 +88,7 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
public List<String> getServiceNames() throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ZipkinServiceTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp =
query(schema,
query(false, schema,
SERVICE_TRAFFIC_TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@ -108,7 +109,7 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
public List<String> getRemoteServiceNames(final String serviceName) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ZipkinServiceRelationTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp =
query(schema,
query(false, schema,
REMOTE_SERVICE_TRAFFIC_TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@ -132,7 +133,7 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
public List<String> getSpanNames(final String serviceName) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ZipkinServiceSpanTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp =
query(schema,
query(false, schema,
SPAN_TRAFFIC_TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@ -153,9 +154,10 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
}
@Override
public List<Span> getTrace(final String traceId) throws IOException {
public List<Span> getTrace(final String traceId, @Nullable final Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp =
query(ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS,
query(isColdStage, ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
@ -197,13 +199,14 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
scrollEndTime = spans.get(spans.size() - 1).getTimestampMillis();
}
return getTraces(traceIds);
return getTraces(traceIds, duration);
}
private List<ZipkinSpanRecord> getSpans(final QueryRequest request,
Duration duration,
long scrollEndTime,
int limit) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
final long startTimeMillis = duration.getStartTimestamp();
TimestampRange tsRange = null;
if (startTimeMillis > 0 && scrollEndTime > 0) {
@ -245,7 +248,7 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
query.setLimit(limit);
}
};
StreamQueryResponse resp = query(ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, tsRange, queryBuilder);
StreamQueryResponse resp = query(isColdStage, ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, tsRange, queryBuilder);
List<ZipkinSpanRecord> spans = new ArrayList<>(); //needs to keep order here
for (final RowEntity rowEntity : resp.getElements()) {
ZipkinSpanRecord spanRecord = new ZipkinSpanRecord.Builder().storage2Entity(
@ -256,13 +259,14 @@ public class BanyanDBZipkinQueryDAO extends AbstractBanyanDBDAO implements IZipk
}
@Override
public List<List<Span>> getTraces(final Set<String> traceIds) throws IOException {
public List<List<Span>> getTraces(final Set<String> traceIds, @Nullable final Duration duration) throws IOException {
if (CollectionUtils.isEmpty(traceIds)) {
return Collections.EMPTY_LIST;
}
final boolean isColdStage = duration != null && duration.isColdStage();
List<AbstractCriteria> conditions = new ArrayList<>(traceIds.size());
StreamQueryResponse resp =
queryDebuggable(ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, null,
queryDebuggable(isColdStage, ZipkinSpanRecord.INDEX_NAME, TRACE_TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override

View File

@ -53,7 +53,7 @@ import java.util.stream.Collectors;
@Override
public List<EBPFProfilingSchedule> querySchedules(String taskId) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(EBPFProfilingScheduleRecord.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@Override

View File

@ -30,10 +30,8 @@ import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
import org.apache.skywalking.banyandb.v1.client.PairQueryCondition;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.DownSampling;
import org.apache.skywalking.oap.server.core.analysis.Layer;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.query.PaginationUtils;
import org.apache.skywalking.oap.server.core.query.enumeration.Order;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -63,16 +61,9 @@ public class BanyanDBEventQueryDAO extends AbstractBanyanDBDAO implements IEvent
public Events queryEvents(EventQueryCondition condition) throws Exception {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(Event.INDEX_NAME, DownSampling.Minute);
final Duration time = condition.getTime();
TimestampRange tsRange = null;
if (time != null) {
long startTB = time.getStartTimeBucketInSec();
long endTB = time.getEndTimeBucketInSec();
if (startTB > 0 && endTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
}
MeasureQueryResponse resp = query(schema, TAGS,
Collections.emptySet(), tsRange, buildQuery(Collections.singletonList(condition)));
boolean isColdStage = time != null && time.isColdStage();
MeasureQueryResponse resp = query(isColdStage, schema, TAGS,
Collections.emptySet(), getTimestampRange(time), buildQuery(Collections.singletonList(condition)));
Events events = new Events();
if (resp.size() == 0) {
return events;
@ -88,16 +79,9 @@ public class BanyanDBEventQueryDAO extends AbstractBanyanDBDAO implements IEvent
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(Event.INDEX_NAME, DownSampling.Minute);
// Duration should be same for all conditions
final Duration time = conditionList.get(0).getTime();
TimestampRange tsRange = null;
if (time != null) {
long startTB = time.getStartTimeBucketInSec();
long endTB = time.getEndTimeBucketInSec();
if (startTB > 0 && endTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
}
MeasureQueryResponse resp = query(schema, TAGS,
Collections.emptySet(), tsRange, buildQuery(conditionList));
boolean isColdStage = time != null && time.isColdStage();
MeasureQueryResponse resp = query(isColdStage, schema, TAGS,
Collections.emptySet(), getTimestampRange(time), buildQuery(conditionList));
Events events = new Events();
if (resp.size() == 0) {
return events;

View File

@ -63,7 +63,7 @@ public class BanyanDBHierarchyQueryDAO extends AbstractBanyanDBDAO implements IH
@Override
public List<ServiceHierarchyRelationTraffic> readAllServiceHierarchyRelations() throws Exception {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ServiceHierarchyRelationTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
SERVICE_HIERARCHY_RELATION_TAGS,
Collections.emptySet(), new QueryBuilder<>() {
@Override
@ -87,7 +87,7 @@ public class BanyanDBHierarchyQueryDAO extends AbstractBanyanDBDAO implements IH
public List<InstanceHierarchyRelationTraffic> readInstanceHierarchyRelations(final String instanceId,
final String layer) throws Exception {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ServiceHierarchyRelationTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
INSTANCE_HIERARCHY_RELATION_TAGS,
Collections.emptySet(), buildInstanceRelationsQuery(instanceId, layer)
);

View File

@ -93,7 +93,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
public List<Service> listServices() throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ServiceTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
SERVICE_TRAFFIC_TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@Override
@ -117,7 +117,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
timestampRange = new TimestampRange(0, duration.getEndTimestamp());
}
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(InstanceTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
INSTANCE_TRAFFIC_TAGS,
Collections.emptySet(),
timestampRange,
@ -145,7 +145,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
public ServiceInstance getInstance(String instanceId) throws IOException {
IDManager.ServiceInstanceID.InstanceIDDefinition id = IDManager.ServiceInstanceID.analysisId(instanceId);
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(InstanceTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
INSTANCE_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -161,7 +161,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public List<ServiceInstance> getInstances(List<String> instanceIds) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(InstanceTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
INSTANCE_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -184,7 +184,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public List<Endpoint> findEndpoint(String keyword, String serviceId, int limit, Duration duration) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(EndpointTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
ENDPOINT_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -220,7 +220,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public List<Process> listProcesses(String serviceId, ProfilingSupportStatus supportStatus, long lastPingStartTimeBucket, long lastPingEndTimeBucket) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -253,7 +253,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
public List<Process> listProcesses(String serviceInstanceId, Duration duration, boolean includeVirtual) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
long lastPingStartTimeBucket = duration.getStartTimeBucket();
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -279,7 +279,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public List<Process> listProcesses(String agentId, long startPingTimeBucket, long endPingTimeBucket) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -304,7 +304,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public long getProcessCount(String serviceId, ProfilingSupportStatus profilingSupportStatus, long lastPingStartTimeBucket, long lastPingEndTimeBucket) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -326,7 +326,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public long getProcessCount(String instanceId) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {
@ -346,7 +346,7 @@ public class BanyanDBMetadataQueryDAO extends AbstractBanyanDBDAO implements IMe
@Override
public Process getProcess(String processId) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ProcessTraffic.INDEX_NAME, DownSampling.Minute);
MeasureQueryResponse resp = query(schema,
MeasureQueryResponse resp = query(false, schema,
PROCESS_TRAFFIC_TAGS,
Collections.emptySet(),
new QueryBuilder<MeasureQuery>() {

View File

@ -125,7 +125,7 @@ public class BanyanDBMetricsDAO extends AbstractBanyanDBDAO implements IMetricsD
}
List<Metrics> metricsInStorage = new ArrayList<>(metrics.size());
MeasureQueryResponse resp = query(schema, schema.getTags(), schema.getFields(), timestampRange, new QueryBuilder<MeasureQuery>() {
MeasureQueryResponse resp = query(false, schema, schema.getTags(), schema.getFields(), timestampRange, new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
seriesIDColumns.entrySet().forEach(entry -> {

View File

@ -29,7 +29,6 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.metrics.DataTable;
import org.apache.skywalking.oap.server.core.analysis.metrics.HistogramMetrics;
@ -49,8 +48,6 @@ import org.apache.skywalking.oap.server.storage.plugin.banyandb.MetadataRegistry
import org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractBanyanDBDAO;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.util.ByteUtil;
import static java.util.Objects.nonNull;
@Slf4j
public class BanyanDBMetricsQueryDAO extends AbstractBanyanDBDAO implements IMetricsQueryDAO {
public BanyanDBMetricsQueryDAO(BanyanDBStorageClient client) {
@ -133,22 +130,13 @@ public class BanyanDBMetricsQueryDAO extends AbstractBanyanDBDAO implements IMet
final String valueColumnName,
final List<KeyValue> labels,
final Duration duration) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
TimestampRange timestampRange = null;
if (startTB > 0 && endTB > 0) {
timestampRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
final boolean isColdStage = duration != null && duration.isColdStage();
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(metricsName, duration.getStep());
if (schema == null) {
throw new IOException("schema is not registered");
}
MeasureQueryResponse resp = query(
schema, ImmutableSet.of(Metrics.ENTITY_ID), ImmutableSet.of(valueColumnName), timestampRange,
isColdStage, schema, ImmutableSet.of(Metrics.ENTITY_ID), ImmutableSet.of(valueColumnName), getTimestampRange(duration),
new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
@ -214,10 +202,9 @@ public class BanyanDBMetricsQueryDAO extends AbstractBanyanDBDAO implements IMet
}
private Map<Long, DataPoint> queryByEntityID(MetadataRegistry.Schema schema, String valueColumnName, Duration duration, String entityID) throws IOException {
TimestampRange timestampRange = new TimestampRange(duration.getStartTimestamp(), duration.getEndTimestamp());
final boolean isColdStage = duration != null && duration.isColdStage();
Map<Long, DataPoint> map = new HashMap<>();
MeasureQueryResponse resp = queryDebuggable(schema, ImmutableSet.of(Metrics.ENTITY_ID), ImmutableSet.of(valueColumnName), timestampRange, new QueryBuilder<MeasureQuery>() {
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema, ImmutableSet.of(Metrics.ENTITY_ID), ImmutableSet.of(valueColumnName), getTimestampRange(duration), new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
query.and(eq(Metrics.ENTITY_ID, entityID));

View File

@ -65,6 +65,7 @@ public class BanyanDBNetworkAddressAliasDAO extends AbstractBanyanDBDAO implemen
public List<NetworkAddressAlias> loadLastUpdate(long timeBucket) {
try {
MeasureQueryResponse resp = query(
false,
getSchema(),
TAGS,
Collections.emptySet(),

View File

@ -47,7 +47,7 @@ public class BanyanDBServiceLabelDAO extends AbstractBanyanDBDAO implements ISer
@Override
public List<String> queryAllLabels(String serviceId) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(ServiceLabelRecord.INDEX_NAME, DownSampling.Minute);
return query(schema, TAGS,
return query(false, schema, TAGS,
Collections.emptySet(), new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(final MeasureQuery query) {

View File

@ -53,6 +53,7 @@ public class BanyanDBTagAutocompleteQueryDAO extends AbstractBanyanDBDAO impleme
@Override
public Set<String> queryTagAutocompleteKeys(TagType tagType, int limit, Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(TagAutocompleteData.INDEX_NAME, DownSampling.Minute);
long startMinTB = 0;
long endMinTB = 0;
@ -68,10 +69,10 @@ public class BanyanDBTagAutocompleteQueryDAO extends AbstractBanyanDBDAO impleme
if (startTB > 0 && endTB > 0) {
range = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
MeasureQueryResponse resp = query(schema,
TAGS_KEY, Collections.emptySet(),
range,
new QueryBuilder<MeasureQuery>() {
MeasureQueryResponse resp = query(isColdStage, schema,
TAGS_KEY, Collections.emptySet(),
range,
new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
query.groupBy(ImmutableSet.of(TagAutocompleteData.TAG_KEY));
@ -94,6 +95,7 @@ public class BanyanDBTagAutocompleteQueryDAO extends AbstractBanyanDBDAO impleme
@Override
public Set<String> queryTagAutocompleteValues(TagType tagType, String tagKey, int limit, Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(TagAutocompleteData.INDEX_NAME, DownSampling.Minute);
long startMinTB = 0;
long endMinTB = 0;
@ -109,10 +111,10 @@ public class BanyanDBTagAutocompleteQueryDAO extends AbstractBanyanDBDAO impleme
if (startTB > 0 && endTB > 0) {
range = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
MeasureQueryResponse resp = query(schema,
TAGS_KV, Collections.emptySet(),
range,
new QueryBuilder<MeasureQuery>() {
MeasureQueryResponse resp = query(isColdStage, schema,
TAGS_KV, Collections.emptySet(),
range,
new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
query.groupBy(ImmutableSet.of(TagAutocompleteData.TAG_VALUE));

View File

@ -32,10 +32,8 @@ import org.apache.skywalking.banyandb.v1.client.AbstractCriteria;
import org.apache.skywalking.banyandb.v1.client.DataPoint;
import org.apache.skywalking.banyandb.v1.client.MeasureQuery;
import org.apache.skywalking.banyandb.v1.client.MeasureQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.UnexpectedException;
import org.apache.skywalking.oap.server.core.analysis.DownSampling;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.manual.relation.endpoint.EndpointRelationServerSideMetrics;
import org.apache.skywalking.oap.server.core.analysis.manual.relation.instance.ServiceInstanceRelationClientSideMetrics;
import org.apache.skywalking.oap.server.core.analysis.manual.relation.instance.ServiceInstanceRelationServerSideMetrics;
@ -54,8 +52,6 @@ import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageC
import org.apache.skywalking.oap.server.storage.plugin.banyandb.MetadataRegistry;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.stream.AbstractBanyanDBDAO;
import static java.util.Objects.nonNull;
public class BanyanDBTopologyQueryDAO extends AbstractBanyanDBDAO implements ITopologyQueryDAO {
public BanyanDBTopologyQueryDAO(final BanyanDBStorageClient client) {
@ -110,25 +106,16 @@ public class BanyanDBTopologyQueryDAO extends AbstractBanyanDBDAO implements ITo
List<Call.CallDetail> queryServiceRelation(Duration duration,
QueryBuilder<MeasureQuery> queryBuilder,
DetectPoint detectPoint) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
TimestampRange timestampRange = null;
if (startTB > 0 && endTB > 0) {
timestampRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
final boolean isColdStage = duration != null && duration.isColdStage();
final String modelName = detectPoint == DetectPoint.SERVER ? ServiceRelationServerSideMetrics.INDEX_NAME :
ServiceRelationClientSideMetrics.INDEX_NAME;
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(modelName, duration.getStep());
MeasureQueryResponse resp = queryDebuggable(schema,
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema,
ImmutableSet.of(
ServiceRelationClientSideMetrics.COMPONENT_IDS,
Metrics.ENTITY_ID
),
Collections.emptySet(), timestampRange, queryBuilder
Collections.emptySet(), getTimestampRange(duration), queryBuilder
);
if (resp.size() == 0) {
return Collections.emptyList();
@ -193,24 +180,15 @@ public class BanyanDBTopologyQueryDAO extends AbstractBanyanDBDAO implements ITo
List<Call.CallDetail> queryInstanceRelation(Duration duration,
QueryBuilder<MeasureQuery> queryBuilder,
DetectPoint detectPoint) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
TimestampRange timestampRange = null;
if (startTB > 0 && endTB > 0) {
timestampRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
final boolean isColdStage = duration != null && duration.isColdStage();
final String modelName = detectPoint == DetectPoint.SERVER ? ServiceInstanceRelationServerSideMetrics.INDEX_NAME :
ServiceInstanceRelationClientSideMetrics.INDEX_NAME;
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(modelName, duration.getStep());
MeasureQueryResponse resp = queryDebuggable(schema,
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema,
ImmutableSet.of(
Metrics.ENTITY_ID
),
Collections.emptySet(), timestampRange, queryBuilder
Collections.emptySet(), getTimestampRange(duration), queryBuilder
);
if (resp.size() == 0) {
return Collections.emptyList();
@ -259,22 +237,13 @@ public class BanyanDBTopologyQueryDAO extends AbstractBanyanDBDAO implements ITo
List<Call.CallDetail> queryEndpointRelation(Duration duration,
QueryBuilder<MeasureQuery> queryBuilder,
DetectPoint detectPoint) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
TimestampRange timestampRange = null;
if (startTB > 0 && endTB > 0) {
timestampRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
final boolean isColdStage = duration != null && duration.isColdStage();
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(EndpointRelationServerSideMetrics.INDEX_NAME, duration.getStep());
MeasureQueryResponse resp = queryDebuggable(schema,
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema,
ImmutableSet.of(
Metrics.ENTITY_ID
),
Collections.emptySet(), timestampRange, queryBuilder
Collections.emptySet(), getTimestampRange(duration), queryBuilder
);
if (resp.size() == 0) {
return Collections.emptyList();
@ -292,23 +261,14 @@ public class BanyanDBTopologyQueryDAO extends AbstractBanyanDBDAO implements ITo
List<Call.CallDetail> queryProcessRelation(Duration duration,
String serviceInstanceId,
DetectPoint detectPoint) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
TimestampRange timestampRange = null;
if (startTB > 0 && endTB > 0) {
timestampRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
final boolean isColdStage = duration != null && duration.isColdStage();
final String modelName = detectPoint == DetectPoint.SERVER ? ProcessRelationServerSideMetrics.INDEX_NAME :
ProcessRelationClientSideMetrics.INDEX_NAME;
// process relation only has minute data
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findMetadata(modelName, DownSampling.Minute);
MeasureQueryResponse resp = queryDebuggable(schema,
MeasureQueryResponse resp = queryDebuggable(isColdStage, schema,
ImmutableSet.of(Metrics.ENTITY_ID, ProcessRelationClientSideMetrics.COMPONENT_ID),
Collections.emptySet(), timestampRange, new QueryBuilder<MeasureQuery>() {
Collections.emptySet(), getTimestampRange(duration), new QueryBuilder<MeasureQuery>() {
@Override
protected void apply(MeasureQuery query) {
query.and(eq(ProcessRelationServerSideMetrics.SERVICE_INSTANCE_ID, serviceInstanceId));

View File

@ -20,6 +20,7 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.gson.Gson;
import java.util.Objects;
import javax.annotation.Nullable;
import org.apache.skywalking.banyandb.model.v1.BanyandbModel;
import org.apache.skywalking.banyandb.v1.client.AbstractCriteria;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
@ -36,12 +37,14 @@ import org.apache.skywalking.banyandb.v1.client.TopNQuery;
import org.apache.skywalking.banyandb.v1.client.TopNQueryResponse;
import org.apache.skywalking.banyandb.v1.client.Trace;
import org.apache.skywalking.oap.server.core.query.input.AttrCondition;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.query.type.KeyValue;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingSpan;
import org.apache.skywalking.oap.server.core.query.type.debugging.DebuggingTraceContext;
import org.apache.skywalking.oap.server.core.storage.AbstractDAO;
import org.apache.skywalking.oap.server.library.util.CollectionUtils;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageConfig;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.MetadataRegistry;
import java.io.IOException;
import java.time.Instant;
@ -63,11 +66,17 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
super(client);
}
protected StreamQueryResponse query(String streamModelName, Set<String> tags, QueryBuilder<StreamQuery> builder) throws IOException {
return this.query(streamModelName, tags, null, builder);
protected StreamQueryResponse query(boolean isColdStage,
String streamModelName,
Set<String> tags,
QueryBuilder<StreamQuery> builder) throws IOException {
return this.query(isColdStage, streamModelName, tags, null, builder);
}
protected StreamQueryResponse query(String streamModelName, Set<String> tags, TimestampRange timestampRange,
protected StreamQueryResponse query(boolean isColdStage,
String streamModelName,
Set<String> tags,
TimestampRange timestampRange,
QueryBuilder<StreamQuery> builder) throws IOException {
MetadataRegistry.Schema schema = MetadataRegistry.INSTANCE.findRecordMetadata(streamModelName);
if (schema == null) {
@ -79,6 +88,9 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
} else {
query = new StreamQuery(List.of(schema.getMetadata().getGroup()), schema.getMetadata().name(), timestampRange, tags);
}
if (isColdStage) {
query.setStages(Set.of(BanyanDBStorageConfig.StageName.cold.name()));
}
builder.apply(query);
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
@ -88,7 +100,10 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
return getClient().query(query);
}
protected StreamQueryResponse queryDebuggable(String modelName, Set<String> tags, TimestampRange timestampRange,
protected StreamQueryResponse queryDebuggable(boolean isColdStage,
String modelName,
Set<String> tags,
TimestampRange timestampRange,
QueryBuilder<StreamQuery> queryBuilder) throws IOException {
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
DebuggingSpan span = null;
@ -105,10 +120,12 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
.append(", Tags: ")
.append(tags)
.append(", TimestampRange: ")
.append(timestampRange);
.append(timestampRange)
.append(", Is cold data query: ")
.append(isColdStage);
span.setMsg(builder.toString());
}
StreamQueryResponse response = query(modelName, tags, timestampRange, queryBuilder);
StreamQueryResponse response = query(isColdStage, modelName, tags, timestampRange, queryBuilder);
if (traceContext != null && traceContext.isDumpStorageRsp()) {
builder.append("\n").append(" Response: ").append(new Gson().toJson(response.getElements()));
span.setMsg(builder.toString());
@ -138,7 +155,8 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
return topNQuery(schema, timestampRange, number, AbstractQuery.Sort.ASC, additionalConditions, attributes);
}
protected TopNQueryResponse topNQueryDebuggable(MetadataRegistry.Schema schema,
protected TopNQueryResponse topNQueryDebuggable(boolean isColdStage,
MetadataRegistry.Schema schema,
TimestampRange timestampRange,
int number,
AbstractQuery.Sort sort,
@ -162,7 +180,9 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
.append(", AdditionalConditions: ")
.append(additionalConditions)
.append(", Attributes: ")
.append(attributes);
.append(attributes)
.append(", Is cold data query: ")
.append(isColdStage);
span.setMsg(builder.toString());
}
TopNQueryResponse response = topNQuery(schema, timestampRange, number, sort, additionalConditions, attributes);
@ -209,7 +229,8 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
return getClient().query(q);
}
protected MeasureQueryResponse queryDebuggable(MetadataRegistry.Schema schema,
protected MeasureQueryResponse queryDebuggable(boolean isColdStage,
MetadataRegistry.Schema schema,
Set<String> tags,
Set<String> fields,
TimestampRange timestampRange,
@ -228,10 +249,12 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
.append(", Fields: ")
.append(fields)
.append(", TimestampRange: ")
.append(timestampRange);
.append(timestampRange)
.append(", Is cold data query: ")
.append(isColdStage);
span.setMsg(builder.toString());
}
MeasureQueryResponse response = query(schema, tags, fields, timestampRange, queryBuilder);
MeasureQueryResponse response = query(isColdStage, schema, tags, fields, timestampRange, queryBuilder);
if (traceContext != null && traceContext.isDumpStorageRsp()) {
builder.append("\n").append(" Response: ").append(new Gson().toJson(response.getDataPoints()));
span.setMsg(builder.toString());
@ -245,15 +268,20 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
}
}
protected MeasureQueryResponse query(MetadataRegistry.Schema schema,
protected MeasureQueryResponse query(boolean isColdStage,
MetadataRegistry.Schema schema,
Set<String> tags,
Set<String> fields,
QueryBuilder<MeasureQuery> builder) throws IOException {
return query(schema, tags, fields, null, builder);
return query(isColdStage, schema, tags, fields, null, builder);
}
protected MeasureQueryResponse query(MetadataRegistry.Schema schema, Set<String> tags, Set<String> fields,
TimestampRange timestampRange, QueryBuilder<MeasureQuery> builder) throws IOException {
protected MeasureQueryResponse query(boolean isColdStage,
MetadataRegistry.Schema schema,
Set<String> tags,
Set<String> fields,
TimestampRange timestampRange,
QueryBuilder<MeasureQuery> builder) throws IOException {
if (schema == null) {
throw new IllegalArgumentException("measure is not registered");
}
@ -263,6 +291,9 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
} else {
query = new MeasureQuery(List.of(schema.getMetadata().getGroup()), schema.getMetadata().name(), timestampRange, tags, fields);
}
if (isColdStage) {
query.setStages(Set.of(BanyanDBStorageConfig.StageName.cold.name()));
}
builder.apply(query);
DebuggingTraceContext traceContext = DebuggingTraceContext.TRACE_CONTEXT.get();
@ -378,4 +409,20 @@ public abstract class AbstractBanyanDBDAO extends AbstractDAO<BanyanDBStorageCli
Or::create);
}
}
protected TimestampRange getTimestampRange(@Nullable Duration duration) {
long startTimeMillis = 0;
long endTimeMillis = 0;
if (duration != null) {
startTimeMillis = duration.getStartTimestamp();
endTimeMillis = duration.getEndTimestamp();
}
TimestampRange tsRange = null;
if (startTimeMillis > 0 && endTimeMillis > 0) {
tsRange = new TimestampRange(startTimeMillis, endTimeMillis);
}
return tsRange;
}
}

View File

@ -23,9 +23,7 @@ import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.alarm.AlarmRecord;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.query.type.AlarmMessage;
@ -57,15 +55,9 @@ public class BanyanDBAlarmQueryDAO extends AbstractBanyanDBDAO implements IAlarm
@Override
public Alarms getAlarm(Integer scopeId, String keyword, int limit, int from, Duration duration, List<Tag> tags) throws IOException {
long startTB = duration.getStartTimeBucketInSec();
long endTB = duration.getEndTimeBucketInSec();
TimestampRange tsRange = null;
if (startTB > 0 && endTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
StreamQueryResponse resp = query(AlarmRecord.INDEX_NAME, TAGS,
tsRange,
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp = query(isColdStage, AlarmRecord.INDEX_NAME, TAGS,
getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {

View File

@ -54,7 +54,7 @@ public class BanyanDBAsyncProfilerTaskLogQueryDAO extends AbstractBanyanDBDAO im
@Override
public List<AsyncProfilerTaskLog> getTaskLogList() throws IOException {
StreamQueryResponse resp = query(AsyncProfilerTaskLogRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, AsyncProfilerTaskLogRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {

View File

@ -70,7 +70,7 @@ public class BanyanDBAsyncProfilerTaskQueryDAO extends AbstractBanyanDBDAO imple
if (endTimeBucket != null) {
endTS = TimeBucket.getTimestamp(endTimeBucket);
}
StreamQueryResponse resp = query(AsyncProfilerTaskRecord.INDEX_NAME, TAGS, new TimestampRange(startTS, endTS),
StreamQueryResponse resp = query(false, AsyncProfilerTaskRecord.INDEX_NAME, TAGS, new TimestampRange(startTS, endTS),
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
@ -96,7 +96,7 @@ public class BanyanDBAsyncProfilerTaskQueryDAO extends AbstractBanyanDBDAO imple
@Override
public AsyncProfilerTask getById(String id) throws IOException {
StreamQueryResponse resp = query(AsyncProfilerTaskRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, AsyncProfilerTaskRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {

View File

@ -22,8 +22,6 @@ import com.google.common.collect.ImmutableSet;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.browser.manual.errorlog.BrowserErrorLogRecord;
import org.apache.skywalking.oap.server.core.browser.source.BrowserErrorCategory;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -32,7 +30,6 @@ import org.apache.skywalking.oap.server.core.query.type.BrowserErrorLogs;
import org.apache.skywalking.oap.server.core.storage.query.IBrowserLogQueryDAO;
import org.apache.skywalking.oap.server.library.util.StringUtil;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
import java.io.IOException;
import java.util.Objects;
import java.util.Set;
@ -53,15 +50,9 @@ public class BanyanDBBrowserLogQueryDAO extends AbstractBanyanDBDAO implements I
public BrowserErrorLogs queryBrowserErrorLogs(String serviceId, String serviceVersionId, String pagePathId,
BrowserErrorCategory category, Duration duration,
int limit, int from) throws IOException {
long startSecondTB = duration.getStartTimeBucketInSec();
long endSecondTB = duration.getEndTimeBucketInSec();
TimestampRange tsRange = null;
if (startSecondTB > 0 && endSecondTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startSecondTB), TimeBucket.getTimestamp(endSecondTB));
}
StreamQueryResponse resp = query(BrowserErrorLogRecord.INDEX_NAME, TAGS,
tsRange, new QueryBuilder<StreamQuery>() {
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp = query(isColdStage, BrowserErrorLogRecord.INDEX_NAME, TAGS,
getTimestampRange(duration), new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {
if (StringUtil.isNotEmpty(serviceId)) {

View File

@ -51,7 +51,7 @@ public class BanyanDBEBPFProfilingDataDAO extends AbstractBanyanDBDAO implements
public List<EBPFProfilingDataRecord> queryData(List<String> scheduleIdList, long beginTime, long endTime) throws IOException {
List<EBPFProfilingDataRecord> records = new ArrayList<>();
for (final String scheduleId : scheduleIdList) {
StreamQueryResponse resp = query(EBPFProfilingDataRecord.INDEX_NAME,
StreamQueryResponse resp = query(false, EBPFProfilingDataRecord.INDEX_NAME,
TAGS,
new QueryBuilder<StreamQuery>() {
@Override

View File

@ -62,7 +62,7 @@ public class BanyanDBEBPFProfilingTaskDAO extends AbstractBanyanDBDAO implements
long taskStartTime, long latestUpdateTime) throws IOException {
List<EBPFProfilingTaskRecord> tasks = new ArrayList<>();
for (final String serviceId : serviceIdList) {
StreamQueryResponse resp = query(EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
@ -85,7 +85,7 @@ public class BanyanDBEBPFProfilingTaskDAO extends AbstractBanyanDBDAO implements
EBPFProfilingTriggerType triggerType, long taskStartTime, long latestUpdateTime) throws IOException {
List<EBPFProfilingTaskRecord> tasks = new ArrayList<>();
for (final EBPFProfilingTargetType targetType : targetTypes) {
StreamQueryResponse resp = query(EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
@ -111,7 +111,7 @@ public class BanyanDBEBPFProfilingTaskDAO extends AbstractBanyanDBDAO implements
@Override
public List<EBPFProfilingTaskRecord> getTaskRecord(String id) throws IOException {
StreamQueryResponse resp = query(EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, EBPFProfilingTaskRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {

View File

@ -53,7 +53,7 @@ public class BanyanDBJFRDataQueryDAO extends AbstractBanyanDBDAO implements IJFR
if (StringUtil.isBlank(taskId) || StringUtil.isBlank(eventType)) {
return new ArrayList<>();
}
StreamQueryResponse resp = query(JFRProfilingDataRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, JFRProfilingDataRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {

View File

@ -22,9 +22,7 @@ import com.google.common.collect.ImmutableSet;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.IDManager;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.manual.log.AbstractLogRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
@ -38,15 +36,12 @@ import org.apache.skywalking.oap.server.core.storage.query.ILogQueryDAO;
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.banyandb.BanyanDBStorageClient;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
import java.util.Set;
import static java.util.Objects.nonNull;
/**
* {@link org.apache.skywalking.oap.server.core.analysis.manual.log.LogRecord} is a stream
*/
@ -72,12 +67,7 @@ public class BanyanDBLogQueryDAO extends AbstractBanyanDBDAO implements ILogQuer
TraceScopeCondition relatedTrace, Order queryOrder, int from, int limit,
Duration duration, List<Tag> tags, List<String> keywordsOfContent,
List<String> excludingKeywordsOfContent) throws IOException {
long startTB = 0;
long endTB = 0;
if (nonNull(duration)) {
startTB = duration.getStartTimeBucketInSec();
endTB = duration.getEndTimeBucketInSec();
}
final boolean isColdStage = duration != null && duration.isColdStage();
final QueryBuilder<StreamQuery> query = new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {
@ -113,12 +103,7 @@ public class BanyanDBLogQueryDAO extends AbstractBanyanDBDAO implements ILogQuer
}
};
TimestampRange tsRange = null;
if (startTB > 0 && endTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startTB), TimeBucket.getTimestamp(endTB));
}
StreamQueryResponse resp = queryDebuggable(LogRecord.INDEX_NAME, TAGS, tsRange, query);
StreamQueryResponse resp = queryDebuggable(isColdStage, LogRecord.INDEX_NAME, TAGS, getTimestampRange(duration), query);
Logs logs = new Logs();

View File

@ -50,7 +50,7 @@ public class BanyanDBProfileTaskLogQueryDAO extends AbstractBanyanDBDAO implemen
@Override
public List<ProfileTaskLog> getTaskLogList() throws IOException {
StreamQueryResponse resp = query(ProfileTaskLogRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, ProfileTaskLogRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {

View File

@ -67,7 +67,7 @@ public class BanyanDBProfileTaskQueryDAO extends AbstractBanyanDBDAO implements
if (endTimeBucket != null) {
endTS = TimeBucket.getTimestamp(endTimeBucket);
}
StreamQueryResponse resp = query(ProfileTaskRecord.INDEX_NAME, TAGS, new TimestampRange(startTS, endTS),
StreamQueryResponse resp = query(false, ProfileTaskRecord.INDEX_NAME, TAGS, new TimestampRange(startTS, endTS),
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
@ -101,7 +101,7 @@ public class BanyanDBProfileTaskQueryDAO extends AbstractBanyanDBDAO implements
@Override
public ProfileTask getById(String id) throws IOException {
StreamQueryResponse resp = query(ProfileTaskRecord.INDEX_NAME, TAGS,
StreamQueryResponse resp = query(false, ProfileTaskRecord.INDEX_NAME, TAGS,
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {

View File

@ -79,7 +79,7 @@ public class BanyanDBProfileThreadSnapshotQueryDAO extends AbstractBanyanDBDAO i
@Override
public List<String> queryProfiledSegmentIdList(String taskId) throws IOException {
StreamQueryResponse resp = query(ProfileThreadSnapshotRecord.INDEX_NAME,
StreamQueryResponse resp = query(false, ProfileThreadSnapshotRecord.INDEX_NAME,
TAGS_BASIC,
new QueryBuilder<StreamQuery>() {
@Override
@ -115,7 +115,7 @@ public class BanyanDBProfileThreadSnapshotQueryDAO extends AbstractBanyanDBDAO i
@Override
public List<ProfileThreadSnapshotRecord> queryRecords(String segmentId, int minSequence, int maxSequence) throws IOException {
StreamQueryResponse resp = query(ProfileThreadSnapshotRecord.INDEX_NAME,
StreamQueryResponse resp = query(false, ProfileThreadSnapshotRecord.INDEX_NAME,
TAGS_ALL,
new QueryBuilder<StreamQuery>() {
@Override
@ -136,7 +136,7 @@ public class BanyanDBProfileThreadSnapshotQueryDAO extends AbstractBanyanDBDAO i
}
private int querySequenceWithAgg(AggType aggType, String segmentId, long start, long end) throws IOException {
StreamQueryResponse resp = query(ProfileThreadSnapshotRecord.INDEX_NAME,
StreamQueryResponse resp = query(false, ProfileThreadSnapshotRecord.INDEX_NAME,
TAGS_ALL, new TimestampRange(start, end),
new QueryBuilder<StreamQuery>() {
@Override

View File

@ -19,12 +19,14 @@
package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.common.collect.ImmutableSet;
import javax.annotation.Nullable;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventTraceType;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.storage.query.ISpanAttachedEventQueryDAO;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBConverter;
import org.apache.skywalking.oap.server.storage.plugin.banyandb.BanyanDBStorageClient;
@ -54,8 +56,11 @@ public class BanyanDBSpanAttachedEventQueryDAO extends AbstractBanyanDBDAO imple
}
@Override
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds) throws IOException {
final StreamQueryResponse resp = queryDebuggable(SpanAttachedEventRecord.INDEX_NAME, TAGS, null, new QueryBuilder<StreamQuery>() {
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
final StreamQueryResponse resp = queryDebuggable(
isColdStage, SpanAttachedEventRecord.INDEX_NAME, TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
protected void apply(StreamQuery query) {
query.and(in(SpanAttachedEventRecord.RELATED_TRACE_ID, traceIds));

View File

@ -20,14 +20,13 @@ package org.apache.skywalking.oap.server.storage.plugin.banyandb.stream;
import com.google.common.base.Strings;
import com.google.common.collect.ImmutableSet;
import javax.annotation.Nullable;
import org.apache.skywalking.banyandb.v1.client.AbstractQuery;
import org.apache.skywalking.banyandb.v1.client.Element;
import org.apache.skywalking.banyandb.v1.client.RowEntity;
import org.apache.skywalking.banyandb.v1.client.StreamQuery;
import org.apache.skywalking.banyandb.v1.client.StreamQueryResponse;
import org.apache.skywalking.banyandb.v1.client.TimestampRange;
import org.apache.skywalking.oap.server.core.analysis.IDManager;
import org.apache.skywalking.oap.server.core.analysis.TimeBucket;
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
import org.apache.skywalking.oap.server.core.analysis.manual.segment.SegmentRecord;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -49,8 +48,6 @@ import java.util.Collections;
import java.util.List;
import java.util.Set;
import static java.util.Objects.nonNull;
public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITraceQueryDAO {
private static final Set<String> BASIC_TAGS = ImmutableSet.of(SegmentRecord.TRACE_ID,
SegmentRecord.IS_ERROR,
@ -80,12 +77,7 @@ public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITrace
@Override
public TraceBrief queryBasicTraces(Duration duration, long minDuration, long maxDuration, String serviceId, String serviceInstanceId, String endpointId, String traceId, int limit, int from, TraceState traceState, QueryOrder queryOrder, List<Tag> tags) throws IOException {
long startSecondTB = 0;
long endSecondTB = 0;
if (nonNull(duration)) {
startSecondTB = duration.getStartTimeBucketInSec();
endSecondTB = duration.getEndTimeBucketInSec();
}
final boolean isColdStage = duration != null && duration.isColdStage();
final QueryBuilder<StreamQuery> q = new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {
@ -145,15 +137,9 @@ public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITrace
}
};
TimestampRange tsRange = null;
if (startSecondTB > 0 && endSecondTB > 0) {
tsRange = new TimestampRange(TimeBucket.getTimestamp(startSecondTB), TimeBucket.getTimestamp(endSecondTB));
}
StreamQueryResponse resp = queryDebuggable(SegmentRecord.INDEX_NAME,
BASIC_TAGS,
tsRange, q);
StreamQueryResponse resp = queryDebuggable(isColdStage, SegmentRecord.INDEX_NAME,
BASIC_TAGS,
getTimestampRange(duration), q);
TraceBrief traceBrief = new TraceBrief();
@ -182,9 +168,10 @@ public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITrace
}
@Override
public List<SegmentRecord> queryByTraceId(String traceId) throws IOException {
StreamQueryResponse resp = queryDebuggable(SegmentRecord.INDEX_NAME, TAGS, null,
new QueryBuilder<StreamQuery>() {
public List<SegmentRecord> queryByTraceId(String traceId, @Nullable Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp = queryDebuggable(isColdStage, SegmentRecord.INDEX_NAME, TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {
query.and(eq(SegmentRecord.TRACE_ID, traceId));
@ -195,8 +182,9 @@ public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITrace
}
@Override
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList) throws IOException {
StreamQueryResponse resp = query(SegmentRecord.INDEX_NAME, TAGS,
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList, @Nullable Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp = query(isColdStage, SegmentRecord.INDEX_NAME, TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {
@ -208,8 +196,9 @@ public class BanyanDBTraceQueryDAO extends AbstractBanyanDBDAO implements ITrace
}
@Override
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList) throws IOException {
StreamQueryResponse resp = query(SegmentRecord.INDEX_NAME, TAGS,
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList, @Nullable Duration duration) throws IOException {
final boolean isColdStage = duration != null && duration.isColdStage();
StreamQueryResponse resp = query(isColdStage, SegmentRecord.INDEX_NAME, TAGS, getTimestampRange(duration),
new QueryBuilder<StreamQuery>() {
@Override
public void apply(StreamQuery query) {

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.query;
import javax.annotation.Nullable;
import org.apache.skywalking.library.elasticsearch.requests.search.BoolQueryBuilder;
import org.apache.skywalking.library.elasticsearch.requests.search.Query;
import org.apache.skywalking.library.elasticsearch.requests.search.Search;
@ -27,6 +28,7 @@ import org.apache.skywalking.library.elasticsearch.requests.search.Sort;
import org.apache.skywalking.library.elasticsearch.response.search.SearchHit;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventTraceType;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.storage.query.ISpanAttachedEventQueryDAO;
import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchScroller;
@ -54,7 +56,7 @@ public class SpanAttachedEventEsDAO extends EsDAO implements ISpanAttachedEventQ
}
@Override
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds) throws IOException {
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) throws IOException {
final String index =
IndexController.LogicIndicesRegister.getPhysicalTableName(SpanAttachedEventRecord.INDEX_NAME);
final BoolQueryBuilder query = Query.bool();

View File

@ -23,6 +23,7 @@ import java.io.IOException;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import javax.annotation.Nullable;
import org.apache.skywalking.library.elasticsearch.requests.search.BoolQueryBuilder;
import org.apache.skywalking.library.elasticsearch.requests.search.Query;
import org.apache.skywalking.library.elasticsearch.requests.search.RangeQueryBuilder;
@ -178,7 +179,7 @@ public class TraceQueryEsDAO extends EsDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryByTraceId(String traceId) throws IOException {
public List<SegmentRecord> queryByTraceId(String traceId, @Nullable Duration duration) throws IOException {
final String index =
IndexController.LogicIndicesRegister.getPhysicalTableName(SegmentRecord.INDEX_NAME);
@ -196,7 +197,7 @@ public class TraceQueryEsDAO extends EsDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList) throws IOException {
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList, @Nullable Duration duration) throws IOException {
final String index =
IndexController.LogicIndicesRegister.getPhysicalTableName(SegmentRecord.INDEX_NAME);
@ -212,7 +213,7 @@ public class TraceQueryEsDAO extends EsDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList) throws IOException {
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList, @Nullable Duration duration) throws IOException {
final String index =
IndexController.LogicIndicesRegister.getPhysicalTableName(SegmentRecord.INDEX_NAME);

View File

@ -18,6 +18,7 @@
package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.query.zipkin;
import javax.annotation.Nullable;
import org.apache.skywalking.library.elasticsearch.requests.search.BoolQueryBuilder;
import org.apache.skywalking.library.elasticsearch.requests.search.Query;
import org.apache.skywalking.library.elasticsearch.requests.search.Search;
@ -129,7 +130,7 @@ public class ZipkinQueryEsDAO extends EsDAO implements IZipkinQueryDAO {
}
@Override
public List<Span> getTrace(final String traceId) {
public List<Span> getTrace(final String traceId, @Nullable final Duration duration) {
String index = IndexController.LogicIndicesRegister.getPhysicalTableName(ZipkinSpanRecord.INDEX_NAME);
BoolQueryBuilder query = Query.bool().must(Query.term(ZipkinSpanRecord.TRACE_ID, traceId));
SearchBuilder search = Search.builder().query(query).size(SCROLLING_BATCH_SIZE);
@ -216,11 +217,11 @@ public class ZipkinQueryEsDAO extends EsDAO implements IZipkinQueryDAO {
traceIds.add((String) idBucket.get("key"));
}
}
return getTraces(traceIds);
return getTraces(traceIds, duration);
}
@Override
public List<List<Span>> getTraces(final Set<String> traceIds) {
public List<List<Span>> getTraces(final Set<String> traceIds, @Nullable final Duration duration) {
String index = IndexController.LogicIndicesRegister.getPhysicalTableName(ZipkinSpanRecord.INDEX_NAME);
BoolQueryBuilder query = Query.bool().must(Query.terms(ZipkinSpanRecord.TRACE_ID, new ArrayList<>(traceIds)));
SearchBuilder search = Search.builder().query(query).sort(ZipkinSpanRecord.TIMESTAMP_MILLIS, Sort.Order.DESC)

View File

@ -18,10 +18,12 @@
package org.apache.skywalking.oap.server.storage.plugin.jdbc.common.dao;
import javax.annotation.Nullable;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventTraceType;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.storage.query.ISpanAttachedEventQueryDAO;
import org.apache.skywalking.oap.server.library.client.jdbc.hikaricp.JDBCClient;
import org.apache.skywalking.oap.server.library.util.StringUtil;
@ -44,7 +46,7 @@ public class JDBCSpanAttachedEventQueryDAO implements ISpanAttachedEventQueryDAO
@Override
@SneakyThrows
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds) {
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) {
final var tables = tableHelper.getTablesWithinTTL(SpanAttachedEventRecord.INDEX_NAME);
final var results = new ArrayList<SpanAttachedEventRecord>();

View File

@ -19,6 +19,7 @@
package org.apache.skywalking.oap.server.storage.plugin.jdbc.common.dao;
import com.google.common.base.Strings;
import javax.annotation.Nullable;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import lombok.extern.slf4j.Slf4j;
@ -232,7 +233,7 @@ public class JDBCTraceQueryDAO implements ITraceQueryDAO {
@Override
@SneakyThrows
public List<SegmentRecord> queryByTraceId(String traceId) throws IOException {
public List<SegmentRecord> queryByTraceId(String traceId, @Nullable Duration duration) throws IOException {
final var tables = tableHelper.getTablesWithinTTL(SegmentRecord.INDEX_NAME);
final var segmentRecords = new ArrayList<SegmentRecord>();
@ -253,7 +254,7 @@ public class JDBCTraceQueryDAO implements ITraceQueryDAO {
@SneakyThrows
@Override
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList) throws IOException {
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList, @Nullable Duration duration) throws IOException {
final var tables = tableHelper.getTablesWithinTTL(SegmentRecord.INDEX_NAME);
final var segmentRecords = new ArrayList<SegmentRecord>();
final ArrayList<String> conditions = new ArrayList<>();
@ -278,7 +279,7 @@ public class JDBCTraceQueryDAO implements ITraceQueryDAO {
@SneakyThrows
@Override
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList) throws IOException {
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList, @Nullable Duration duration) throws IOException {
final var tables = tableHelper.getTablesWithinTTL(SegmentRecord.INDEX_NAME);
final var segmentRecords = new ArrayList<SegmentRecord>();
final ArrayList<String> conditions = new ArrayList<>();

View File

@ -21,6 +21,7 @@ package org.apache.skywalking.oap.server.storage.plugin.jdbc.common.dao;
import com.google.gson.Gson;
import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import javax.annotation.Nullable;
import lombok.RequiredArgsConstructor;
import lombok.SneakyThrows;
import org.apache.skywalking.oap.server.core.query.input.Duration;
@ -147,7 +148,7 @@ public class JDBCZipkinQueryDAO implements IZipkinQueryDAO {
@Override
@SneakyThrows
public List<Span> getTrace(final String traceId) {
public List<Span> getTrace(final String traceId, @Nullable final Duration duration) {
final var tables = tableHelper.getTablesWithinTTL(ZipkinSpanRecord.INDEX_NAME);
final var trace = new ArrayList<Span>();
@ -264,12 +265,12 @@ public class JDBCZipkinQueryDAO implements IZipkinQueryDAO {
}, condition.toArray(new Object[0]));
}
return getTraces(traceIds);
return getTraces(traceIds, duration);
}
@Override
@SneakyThrows
public List<List<Span>> getTraces(final Set<String> traceIds) {
public List<List<Span>> getTraces(final Set<String> traceIds, final Duration duration) {
if (CollectionUtils.isEmpty(traceIds)) {
return new ArrayList<>();
}

View File

@ -88,7 +88,7 @@ public class ProfiledBasicInfo {
profiledSegment.setDuration(segment.getLatency());
// query spans
Trace trace = traceQueryService.queryTrace(config.getTraceId());
Trace trace = traceQueryService.queryTrace(config.getTraceId(), null);
List<Span> profiledSegmentSpans = trace.getSpans().stream().filter(s -> Objects.equals(s.getSegmentId(), segmentId)).collect(Collectors.toList());
if (CollectionUtils.isEmpty(profiledSegmentSpans)) {
throw new IllegalArgumentException("Current segment cannot found any span");

View File

@ -21,6 +21,7 @@ package org.apache.skywalking.oap.server.tool.profile.exporter.test;
import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import javax.annotation.Nullable;
import org.apache.skywalking.apm.network.language.agent.v3.SegmentObject;
import org.apache.skywalking.apm.network.language.agent.v3.SpanObject;
import org.apache.skywalking.oap.server.core.analysis.manual.searchtag.Tag;
@ -56,7 +57,7 @@ public class ProfileTraceDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryByTraceId(String traceId) throws IOException {
public List<SegmentRecord> queryByTraceId(String traceId, @Nullable Duration duration) throws IOException {
final ArrayList<SegmentRecord> segments = new ArrayList<>();
final SegmentRecord segment = new SegmentRecord();
segments.add(segment);
@ -80,7 +81,7 @@ public class ProfileTraceDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList) throws IOException {
public List<SegmentRecord> queryBySegmentIdList(List<String> segmentIdList, @Nullable Duration duration) throws IOException {
final ArrayList<SegmentRecord> segments = new ArrayList<>();
final SegmentRecord segment = new SegmentRecord();
segments.add(segment);
@ -104,7 +105,7 @@ public class ProfileTraceDAO implements ITraceQueryDAO {
}
@Override
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList) throws IOException {
public List<SegmentRecord> queryByTraceIdWithInstanceId(List<String> traceIdList, List<String> instanceIdList, @Nullable Duration duration) throws IOException {
return null;
}

View File

@ -18,8 +18,10 @@
package org.apache.skywalking.oap.server.tool.profile.exporter.test;
import javax.annotation.Nullable;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventRecord;
import org.apache.skywalking.oap.server.core.analysis.manual.spanattach.SpanAttachedEventTraceType;
import org.apache.skywalking.oap.server.core.query.input.Duration;
import org.apache.skywalking.oap.server.core.storage.query.ISpanAttachedEventQueryDAO;
import java.io.IOException;
@ -27,7 +29,7 @@ import java.util.List;
public class SpanAttachedEventQueryDAO implements ISpanAttachedEventQueryDAO {
@Override
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds) throws IOException {
public List<SpanAttachedEventRecord> querySpanAttachedEvents(SpanAttachedEventTraceType type, List<String> traceIds, @Nullable Duration duration) throws IOException {
return null;
}
}
}

View File

@ -22,7 +22,7 @@ SW_AGENT_PYTHON_COMMIT=c76a6ec51a478ac91abb20ec8f22a99b8d4d6a58
SW_AGENT_CLIENT_JS_COMMIT=af0565a67d382b683c1dbd94c379b7080db61449
SW_AGENT_CLIENT_JS_TEST_COMMIT=4f1eb1dcdbde3ec4a38534bf01dded4ab5d2f016
SW_KUBERNETES_COMMIT_SHA=6fe5e6f0d3b7686c6be0457733e825ee68cb9b35
SW_ROVER_COMMIT=4c0cb8429a96f190ea30eac1807008d523c749c3
SW_ROVER_COMMIT=738b1a42fe4941e0b4e6f5816403437cf572708f
SW_BANYANDB_COMMIT=458041a561b0acc1f2ed37690df2ce753b791283
SW_AGENT_PHP_COMMIT=3192c553002707d344bd6774cfab5bc61f67a1d3
SW_PREDICTOR_COMMIT=54a0197654a3781a6f73ce35146c712af297c994