diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentDispatcher.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentDispatcher.java index c7b11da78..7f4e34b8e 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentDispatcher.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentDispatcher.java @@ -32,6 +32,7 @@ public class SegmentDispatcher implements SourceDispatcher { segment.setSegmentId(source.getSegmentId()); segment.setTraceId(source.getTraceId()); segment.setServiceId(source.getServiceId()); + segment.setServiceInstanceId(source.getServiceInstanceId()); segment.setEndpointName(source.getEndpointName()); segment.setEndpointId(source.getEndpointId()); segment.setStartTime(source.getStartTime()); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java index 6be26acbc..b98e774e6 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/manual/segment/SegmentRecord.java @@ -40,6 +40,7 @@ public class SegmentRecord extends Record { public static final String SEGMENT_ID = "segment_id"; public static final String TRACE_ID = "trace_id"; public static final String SERVICE_ID = "service_id"; + public static final String SERVICE_INSTANCE_ID = "service_instance_id"; public static final String ENDPOINT_NAME = "endpoint_name"; public static final String ENDPOINT_ID = "endpoint_id"; public static final String START_TIME = "start_time"; @@ -52,6 +53,7 @@ public class SegmentRecord extends Record { @Setter @Getter @Column(columnName = SEGMENT_ID) @IDColumn private String segmentId; @Setter @Getter @Column(columnName = TRACE_ID) @IDColumn private String traceId; @Setter @Getter @Column(columnName = SERVICE_ID) @IDColumn private int serviceId; + @Setter @Getter @Column(columnName = SERVICE_INSTANCE_ID) @IDColumn private int serviceInstanceId; @Setter @Getter @Column(columnName = ENDPOINT_NAME, matchQuery = true) @IDColumn private String endpointName; @Setter @Getter @Column(columnName = ENDPOINT_ID) @IDColumn private int endpointId; @Setter @Getter @Column(columnName = START_TIME) @IDColumn private long startTime; @@ -72,6 +74,7 @@ public class SegmentRecord extends Record { map.put(SEGMENT_ID, storageData.getSegmentId()); map.put(TRACE_ID, storageData.getTraceId()); map.put(SERVICE_ID, storageData.getServiceId()); + map.put(SERVICE_INSTANCE_ID, storageData.getServiceInstanceId()); map.put(ENDPOINT_NAME, storageData.getEndpointName()); map.put(ENDPOINT_ID, storageData.getEndpointId()); map.put(START_TIME, storageData.getStartTime()); @@ -93,6 +96,7 @@ public class SegmentRecord extends Record { record.setSegmentId((String)dbMap.get(SEGMENT_ID)); record.setTraceId((String)dbMap.get(TRACE_ID)); record.setServiceId(((Number)dbMap.get(SERVICE_ID)).intValue()); + record.setServiceInstanceId(((Number)dbMap.get(SERVICE_INSTANCE_ID)).intValue()); record.setEndpointName((String)dbMap.get(ENDPOINT_NAME)); record.setEndpointId(((Number)dbMap.get(ENDPOINT_ID)).intValue()); record.setStartTime(((Number)dbMap.get(START_TIME)).longValue()); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TraceQueryService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TraceQueryService.java index b23a9f442..b73552501 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TraceQueryService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/query/TraceQueryService.java @@ -89,14 +89,14 @@ public class TraceQueryService implements Service { return componentLibraryCatalogService; } - public TraceBrief queryBasicTraces(final int serviceId, final int endpointId, final String traceId, - final String endpointName, - final int minTraceDuration, int maxTraceDuration, final TraceState traceState, final QueryOrder queryOrder, + public TraceBrief queryBasicTraces(final int serviceId, final int serviceInstanceId, final int endpointId, + final String traceId, final String endpointName, final int minTraceDuration, int maxTraceDuration, + final TraceState traceState, final QueryOrder queryOrder, final Pagination paging, final long startTB, final long endTB) throws IOException { PaginationUtils.Page page = PaginationUtils.INSTANCE.exchange(paging); return getTraceQueryDAO().queryBasicTraces(startTB, endTB, minTraceDuration, maxTraceDuration, endpointName, - serviceId, endpointId, traceId, page.getLimit(), page.getFrom(), traceState, queryOrder); + serviceId, serviceInstanceId, endpointId, traceId, page.getLimit(), page.getFrom(), traceState, queryOrder); } public Trace queryTrace(final String traceId) throws IOException { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Segment.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Segment.java index 041f7df1a..5959aeda3 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Segment.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/Segment.java @@ -38,6 +38,7 @@ public class Segment extends Source { @Setter @Getter private String segmentId; @Setter @Getter private String traceId; @Setter @Getter private int serviceId; + @Setter @Getter private int serviceInstanceId; @Setter @Getter private String endpointName; @Setter @Getter private int endpointId; @Setter @Getter private long startTime; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITraceQueryDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITraceQueryDAO.java index 2fb4d8787..fc98d03ba 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITraceQueryDAO.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/query/ITraceQueryDAO.java @@ -30,8 +30,8 @@ import org.apache.skywalking.oap.server.library.module.Service; public interface ITraceQueryDAO extends Service { TraceBrief queryBasicTraces(long startSecondTB, long endSecondTB, long minDuration, - long maxDuration, String endpointName, int serviceId, int endpointId, String traceId, int limit, int from, - TraceState traceState, QueryOrder queryOrder) throws IOException; + long maxDuration, String endpointName, int serviceId, int serviceInstanceId, int endpointId, String traceId, + int limit, int from, TraceState traceState, QueryOrder queryOrder) throws IOException; List queryByTraceId(String traceId) throws IOException; } diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/TraceQuery.java b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/TraceQuery.java index eb9f2926a..9eb6bbc1f 100644 --- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/TraceQuery.java +++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/resolver/TraceQuery.java @@ -68,11 +68,12 @@ public class TraceQuery implements GraphQLQueryResolver { String endpointName = condition.getEndpointName(); int serviceId = StringUtils.isEmpty(condition.getServiceId()) ? 0 : Integer.parseInt(condition.getServiceId()); int endpointId = StringUtils.isEmpty(condition.getEndpointId()) ? 0 : Integer.parseInt(condition.getEndpointId()); + int serviceInstanceId = StringUtils.isEmpty(condition.getServiceInstanceId()) ? 0 : Integer.parseInt(condition.getServiceInstanceId()); TraceState traceState = condition.getTraceState(); QueryOrder queryOrder = condition.getQueryOrder(); Pagination pagination = condition.getPaging(); - return getQueryService().queryBasicTraces(serviceId, endpointId, traceId, endpointName, minDuration, maxDuration, traceState, queryOrder, pagination, startSecondTB, endSecondTB); + return getQueryService().queryBasicTraces(serviceId, serviceInstanceId, endpointId, traceId, endpointName, minDuration, maxDuration, traceState, queryOrder, pagination, startSecondTB, endSecondTB); } public Trace queryTrace(final String traceId) throws IOException { diff --git a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/type/TraceQueryCondition.java b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/type/TraceQueryCondition.java index 6db4e7992..0b315c265 100644 --- a/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/type/TraceQueryCondition.java +++ b/oap-server/server-query-plugin/query-graphql-plugin/src/main/java/org/apache/skywalking/oap/query/graphql/type/TraceQueryCondition.java @@ -25,6 +25,7 @@ import org.apache.skywalking.oap.server.core.query.entity.*; @Setter public class TraceQueryCondition { private String serviceId; + private String serviceInstanceId; private String traceId; private String endpointName; private String endpointId; diff --git a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/segment/SegmentSpanListener.java b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/segment/SegmentSpanListener.java index 644d2d20b..e73def0c4 100644 --- a/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/segment/SegmentSpanListener.java +++ b/oap-server/server-receiver-plugin/skywalking-trace-receiver-plugin/src/main/java/org/apache/skywalking/oap/server/receiver/trace/provider/parser/listener/segment/SegmentSpanListener.java @@ -72,6 +72,7 @@ public class SegmentSpanListener implements FirstSpanListener, EntrySpanListener segment.setSegmentId(segmentCoreInfo.getSegmentId()); segment.setServiceId(segmentCoreInfo.getServiceId()); + segment.setServiceInstanceId(segmentCoreInfo.getServiceInstanceId()); segment.setLatency((int)(segmentCoreInfo.getEndTime() - segmentCoreInfo.getStartTime())); segment.setStartTime(segmentCoreInfo.getStartTime()); segment.setEndTime(segmentCoreInfo.getEndTime()); diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TraceQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TraceQueryEsDAO.java index bffed1fa3..ac5645a6b 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TraceQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/TraceQueryEsDAO.java @@ -44,8 +44,8 @@ public class TraceQueryEsDAO extends EsDAO implements ITraceQueryDAO { @Override public TraceBrief queryBasicTraces(long startSecondTB, long endSecondTB, long minDuration, - long maxDuration, String endpointName, int serviceId, int endpointId, String traceId, int limit, int from, - TraceState traceState, QueryOrder queryOrder) throws IOException { + long maxDuration, String endpointName, int serviceId, int serviceInstanceId, int endpointId, String traceId, + int limit, int from, TraceState traceState, QueryOrder queryOrder) throws IOException { SearchSourceBuilder sourceBuilder = SearchSourceBuilder.searchSource(); BoolQueryBuilder boolQueryBuilder = QueryBuilders.boolQuery(); @@ -72,6 +72,9 @@ public class TraceQueryEsDAO extends EsDAO implements ITraceQueryDAO { if (serviceId != 0) { boolQueryBuilder.must().add(QueryBuilders.termQuery(SegmentRecord.SERVICE_ID, serviceId)); } + if (serviceInstanceId != 0) { + boolQueryBuilder.must().add(QueryBuilders.termQuery(SegmentRecord.SERVICE_INSTANCE_ID, serviceInstanceId)); + } if (endpointId != 0) { boolQueryBuilder.must().add(QueryBuilders.termQuery(SegmentRecord.ENDPOINT_ID, endpointId)); } diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TraceQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TraceQueryDAO.java index 277bd97e6..ed8bf7341 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TraceQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/h2/dao/H2TraceQueryDAO.java @@ -41,8 +41,8 @@ public class H2TraceQueryDAO implements ITraceQueryDAO { @Override public TraceBrief queryBasicTraces(long startSecondTB, long endSecondTB, long minDuration, long maxDuration, - String endpointName, int serviceId, int endpointId, String traceId, int limit, int from, TraceState traceState, - QueryOrder queryOrder) throws IOException { + String endpointName, int serviceId, int serviceInstanceId, int endpointId, String traceId, int limit, int from, + TraceState traceState, QueryOrder queryOrder) throws IOException { StringBuilder sql = new StringBuilder(); List parameters = new ArrayList<>(10); @@ -71,6 +71,10 @@ public class H2TraceQueryDAO implements ITraceQueryDAO { sql.append(" and ").append(SegmentRecord.SERVICE_ID).append(" = ?"); parameters.add(serviceId); } + if (serviceInstanceId != 0) { + sql.append(" and ").append(SegmentRecord.SERVICE_INSTANCE_ID).append(" = ?"); + parameters.add(serviceInstanceId); + } if (endpointId != 0) { sql.append(" and ").append(SegmentRecord.ENDPOINT_ID).append(" = ?"); parameters.add(endpointId); diff --git a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLTraceQueryDAO.java b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLTraceQueryDAO.java index ddd36a596..639fe63ca 100644 --- a/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLTraceQueryDAO.java +++ b/oap-server/server-storage-plugin/storage-jdbc-hikaricp-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/jdbc/mysql/MySQLTraceQueryDAO.java @@ -39,8 +39,8 @@ public class MySQLTraceQueryDAO extends H2TraceQueryDAO { @Override public TraceBrief queryBasicTraces(long startSecondTB, long endSecondTB, long minDuration, long maxDuration, - String endpointName, int serviceId, int endpointId, String traceId, int limit, int from, TraceState traceState, - QueryOrder queryOrder) throws IOException { + String endpointName, int serviceId, int serviceInstanceId, int endpointId, String traceId, int limit, int from, + TraceState traceState, QueryOrder queryOrder) throws IOException { StringBuilder sql = new StringBuilder(); List parameters = new ArrayList<>(10); @@ -69,6 +69,10 @@ public class MySQLTraceQueryDAO extends H2TraceQueryDAO { sql.append(" and ").append(SegmentRecord.SERVICE_ID).append(" = ?"); parameters.add(serviceId); } + if (serviceInstanceId != 0) { + sql.append(" and ").append(SegmentRecord.SERVICE_INSTANCE_ID).append(" = ?"); + parameters.add(serviceInstanceId); + } if (endpointId != 0) { sql.append(" and ").append(SegmentRecord.ENDPOINT_ID).append(" = ?"); parameters.add(endpointId);