Skip to content

Commit 4f5cef6

Browse files
author
Jason Bulicek
authored
Merge pull request #183 from ExpediaDotCom/ConditionallyIndexTagServiceNameIfItExists
If the "service" tag exists in the span, send it to ElasticSearch and
2 parents 8816a55 + d085b65 commit 4f5cef6

4 files changed

Lines changed: 17 additions & 24 deletions

File tree

indexer/src/main/scala/com/expedia/www/haystack/trace/indexer/writers/cassandra/ServiceMetadataStatementBuilder.scala

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import java.time.Instant
2222
import com.datastax.driver.core.Statement
2323
import com.expedia.open.tracing.Span
2424
import com.expedia.www.haystack.trace.commons.clients.cassandra.CassandraSession
25+
import com.expedia.www.haystack.trace.commons.utils.SpanUtils
2526
import com.expedia.www.haystack.trace.indexer.config.entities.ServiceMetadataWriteConfiguration
2627
import org.apache.commons.lang3.StringUtils
2728

@@ -72,8 +73,9 @@ class ServiceMetadataStatementBuilder(cassandra: CassandraSession,
7273
def getAndUpdateServiceMetadata(spans: Iterable[Span]): Seq[Statement] = {
7374
this.synchronized {
7475
spans.foreach(span => {
75-
if (StringUtils.isNotEmpty(span.getServiceName) && StringUtils.isNotEmpty(span.getOperationName)) {
76-
val operationsList = serviceMetadataMap.getOrElseUpdate(span.getServiceName, mutable.Set[String]())
76+
val serviceName = SpanUtils.getEffectiveServiceName(span)
77+
if (StringUtils.isNotEmpty(serviceName) && StringUtils.isNotEmpty(span.getOperationName)) {
78+
val operationsList = serviceMetadataMap.getOrElseUpdate(serviceName, mutable.Set[String]())
7779
if (operationsList.add(span.getOperationName)) {
7880
allOperationCount += 1
7981
}

indexer/src/main/scala/com/expedia/www/haystack/trace/indexer/writers/es/IndexDocumentGenerator.scala

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ import com.expedia.www.haystack.trace.commons.clients.es.document.TraceIndexDoc
2626
import com.expedia.www.haystack.trace.commons.clients.es.document.TraceIndexDoc.{OPERATION_KEY_NAME, SERVICE_KEY_NAME, TagValue}
2727
import com.expedia.www.haystack.trace.commons.config.entities.IndexFieldType.IndexFieldType
2828
import com.expedia.www.haystack.trace.commons.config.entities.{IndexFieldType, WhitelistIndexFieldConfiguration}
29+
import com.expedia.www.haystack.trace.commons.utils.SpanUtils
2930
import org.apache.commons.lang3.StringUtils
3031

3132
import scala.collection.JavaConverters._
@@ -53,11 +54,12 @@ class IndexDocumentGenerator(config: WhitelistIndexFieldConfiguration) extends M
5354
traceStartTime = Math.min(traceStartTime, microsToSecondGranularity(span.getStartTime))
5455
if(span.getParentSpanId == null) rootDuration = span.getDuration
5556

57+
val serviceName = SpanUtils.getEffectiveServiceName(span)
5658
val spanIndexDoc = spanIndices
57-
.find(sp => sp(OPERATION_KEY_NAME).equals(span.getOperationName) && sp(SERVICE_KEY_NAME).equals(span.getServiceName))
59+
.find(sp => sp(OPERATION_KEY_NAME).equals(span.getOperationName) && sp(SERVICE_KEY_NAME).equals(serviceName))
5860
.getOrElse({
5961
val newSpanIndexDoc = mutable.Map[String, Any](
60-
SERVICE_KEY_NAME -> span.getServiceName,
62+
SERVICE_KEY_NAME -> serviceName,
6163
OPERATION_KEY_NAME -> span.getOperationName)
6264
spanIndices.append(newSpanIndexDoc)
6365
newSpanIndexDoc
@@ -68,7 +70,7 @@ class IndexDocumentGenerator(config: WhitelistIndexFieldConfiguration) extends M
6870
}
6971

7072
private def isValidForIndex(span: Span): Boolean = {
71-
StringUtils.isNotEmpty(span.getServiceName) && StringUtils.isNotEmpty(span.getOperationName)
73+
StringUtils.isNotEmpty(SpanUtils.getEffectiveServiceName(span)) && StringUtils.isNotEmpty(span.getOperationName)
7274
}
7375

7476
/**

pom.xml

Lines changed: 0 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -251,21 +251,6 @@
251251
<artifactId>config</artifactId>
252252
</dependency>
253253

254-
<dependency>
255-
<groupId>com.google.protobuf</groupId>
256-
<artifactId>protobuf-java</artifactId>
257-
</dependency>
258-
259-
<dependency>
260-
<groupId>io.grpc</groupId>
261-
<artifactId>grpc-protobuf</artifactId>
262-
</dependency>
263-
264-
<dependency>
265-
<groupId>io.grpc</groupId>
266-
<artifactId>grpc-stub</artifactId>
267-
</dependency>
268-
269254
<dependency>
270255
<groupId>org.scala-lang</groupId>
271256
<artifactId>scala-library</artifactId>

reader/src/main/scala/com/expedia/www/haystack/trace/reader/readers/utils/SpanMerger.scala

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -53,8 +53,12 @@ object SpanMerger {
5353
Span
5454
.newBuilder(serverSpan)
5555
.setParentSpanId(clientSpan.getParentSpanId) // use the parentSpanId of the client span to stitch in the client's trace tree
56-
.addAllTags((clientSpan.getTagsList.asScala ++ auxiliaryCommonTags(clientSpan, serverSpan) ++ auxiliaryClientTags(clientSpan) ++ auxiliaryServerTags(serverSpan)).asJavaCollection)
57-
.clearLogs().addAllLogs((clientSpan.getLogsList.asScala ++ serverSpan.getLogsList.asScala.sortBy(_.getTimestamp)).asJavaCollection)
56+
.addAllTags((clientSpan.getTagsList.asScala
57+
++ auxiliaryCommonTags(clientSpan, serverSpan)
58+
++ auxiliaryClientTags(clientSpan)
59+
++ auxiliaryServerTags(serverSpan)).asJavaCollection)
60+
.clearLogs().addAllLogs((clientSpan.getLogsList.asScala
61+
++ serverSpan.getLogsList.asScala.sortBy(_.getTimestamp)).asJavaCollection)
5862
.build()
5963
}
6064

@@ -108,7 +112,7 @@ object SpanMerger {
108112

109113
private def auxiliaryClientTags(span: Span): List[Tag] =
110114
List(
111-
buildStringTag(AuxiliaryTags.CLIENT_SERVICE_NAME, span.getServiceName),
115+
buildStringTag(AuxiliaryTags.CLIENT_SERVICE_NAME, SpanUtils.getEffectiveServiceName(span)),
112116
buildStringTag(AuxiliaryTags.CLIENT_OPERATION_NAME, span.getOperationName),
113117
buildStringTag(AuxiliaryTags.CLIENT_INFRASTRUCTURE_PROVIDER, extractTagStringValue(span, AuxiliaryTags.INFRASTRUCTURE_PROVIDER)),
114118
buildStringTag(AuxiliaryTags.CLIENT_INFRASTRUCTURE_LOCATION, extractTagStringValue(span, AuxiliaryTags.INFRASTRUCTURE_LOCATION)),
@@ -118,7 +122,7 @@ object SpanMerger {
118122

119123
private def auxiliaryServerTags(span: Span): List[Tag] = {
120124
List(
121-
buildStringTag(AuxiliaryTags.SERVER_SERVICE_NAME, span.getServiceName),
125+
buildStringTag(AuxiliaryTags.SERVER_SERVICE_NAME, SpanUtils.getEffectiveServiceName(span)),
122126
buildStringTag(AuxiliaryTags.SERVER_OPERATION_NAME, span.getOperationName),
123127
buildStringTag(AuxiliaryTags.SERVER_INFRASTRUCTURE_PROVIDER, extractTagStringValue(span, AuxiliaryTags.INFRASTRUCTURE_PROVIDER)),
124128
buildStringTag(AuxiliaryTags.SERVER_INFRASTRUCTURE_LOCATION, extractTagStringValue(span, AuxiliaryTags.INFRASTRUCTURE_LOCATION)),

0 commit comments

Comments
 (0)