Skip to content

Commit dbbde37

Browse files
Kapil Rastogivsen
authored andcommitted
Added endpoint for getting traces details for given traceIds (#170)
* Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds * Added endpoint for getting traces details for given traceIds
1 parent 50547e3 commit dbbde37

11 files changed

Lines changed: 352 additions & 9 deletions

File tree

commons/src/main/scala/com/expedia/www/haystack/trace/commons/clients/cassandra/CassandraSession.scala

Lines changed: 29 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,15 @@ class CassandraSession(config: CassandraConfiguration, factory: ClusterFactory)
5757
.where(QueryBuilder.eq(ID_COLUMN_NAME, bindMarker(ID_COLUMN_NAME))))
5858
}
5959

60+
def selectRawTracesPreparedStmt: PreparedStatement = {
61+
import QueryBuilder.bindMarker
62+
session.prepare(
63+
QueryBuilder
64+
.select()
65+
.from(config.tracesKeyspace.name, config.tracesKeyspace.table)
66+
.where(QueryBuilder.in(ID_COLUMN_NAME, bindMarker(ID_COLUMN_NAME))))
67+
}
68+
6069
def createServiceMetadataInsertPreparedStatement(keyspace: KeyspaceConfiguration): PreparedStatement = {
6170
import QueryBuilder.{bindMarker, ttl}
6271

@@ -93,9 +102,10 @@ class CassandraSession(config: CassandraConfiguration, factory: ClusterFactory)
93102

94103
/**
95104
* create bound statement for writing to cassandra table
96-
* @param serviceName name of service ie primary key in cassandra
97-
* @param operationList name of operation
98-
* @param consistencyLevel consistency level for cassandra write
105+
*
106+
* @param serviceName name of service ie primary key in cassandra
107+
* @param operationList name of operation
108+
* @param consistencyLevel consistency level for cassandra write
99109
* @param insertServiceMetadataStatement prepared statement to use
100110
* @return
101111
*/
@@ -117,9 +127,10 @@ class CassandraSession(config: CassandraConfiguration, factory: ClusterFactory)
117127

118128
/**
119129
* create bound statement for writing to cassandra table
120-
* @param traceId trace id
121-
* @param spanBufferBytes data bytes of spanBuffer that belong to a given trace id
122-
* @param consistencyLevel consistency level for cassandra write
130+
*
131+
* @param traceId trace id
132+
* @param spanBufferBytes data bytes of spanBuffer that belong to a given trace id
133+
* @param consistencyLevel consistency level for cassandra write
123134
* @param insertTraceStatement prepared statement to use
124135
* @return
125136
*/
@@ -136,15 +147,27 @@ class CassandraSession(config: CassandraConfiguration, factory: ClusterFactory)
136147

137148
/**
138149
* create new select statement for retrieving data for traceId
150+
*
139151
* @param traceId trace id
140152
* @return
141153
*/
142154
def newSelectTraceBoundStatement(traceId: String): Statement = {
143155
new BoundStatement(selectTracePreparedStmt).setString(ID_COLUMN_NAME, traceId)
144156
}
145157

158+
/**
159+
* create new select statement for retrieving Raw Traces data for traceIds
160+
*
161+
* @param traceIds list of trace id
162+
* @return statement for select query for traceIds
163+
*/
164+
def newSelectRawTracesBoundStatement(traceIds: List[String]): Statement = {
165+
new BoundStatement(selectRawTracesPreparedStmt).setList(ID_COLUMN_NAME, traceIds.asJava)
166+
}
167+
146168
/**
147169
* executes the statement async and return the resultset future
170+
*
148171
* @param statement prepared statement to be executed
149172
* @return future object of ResultSet
150173
*/

commons/src/test/scala/com/expedia/www/haystack/trace/commons/unit/CassandraSessionSpec.scala

Lines changed: 45 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
package com.expedia.www.haystack.trace.commons.unit
1919

2020
import com.datastax.driver.core._
21-
import com.datastax.driver.core.querybuilder.Insert
21+
import com.datastax.driver.core.querybuilder.{Insert, Select}
2222
import com.expedia.www.haystack.trace.commons.clients.cassandra.{CassandraClusterFactory, CassandraSession}
2323
import com.expedia.www.haystack.trace.commons.config.entities.{CassandraConfiguration, KeyspaceConfiguration, SocketConfiguration}
2424
import org.easymock.EasyMock
@@ -38,7 +38,7 @@ class CassandraSessionSpec extends FunSpec with Matchers with EasyMockSugar {
3838
val keyspaceMetadata = mock[KeyspaceMetadata]
3939
val tableMetadata = mock[TableMetadata]
4040
val insertPrepStatement = mock[PreparedStatement]
41-
val keyspaceConfig = KeyspaceConfiguration(keyspaceName, tableName, 100, None)
41+
val keyspaceConfig = KeyspaceConfiguration(keyspaceName, tableName, 100, None)
4242

4343
val config = CassandraConfiguration(List("cassandra1"),
4444
autoDiscoverEnabled = false,
@@ -69,5 +69,48 @@ class CassandraSessionSpec extends FunSpec with Matchers with EasyMockSugar {
6969
session.close()
7070
}
7171
}
72+
73+
it("should connect to the cassandra cluster and provide prepared statement for select with traces") {
74+
val keyspaceName = "keyspace-1"
75+
val tableName = "table-1"
76+
77+
val factory = mock[CassandraClusterFactory]
78+
val session = mock[Session]
79+
val cluster = mock[Cluster]
80+
val metadata = mock[Metadata]
81+
val keyspaceMetadata = mock[KeyspaceMetadata]
82+
val tableMetadata = mock[TableMetadata]
83+
val selectPrepStatement = mock[PreparedStatement]
84+
val keyspaceConfig = KeyspaceConfiguration(keyspaceName, tableName, 100, None)
85+
86+
val config = CassandraConfiguration(List("cassandra1"),
87+
autoDiscoverEnabled = false,
88+
None,
89+
None,
90+
keyspaceConfig,
91+
SocketConfiguration(10, keepAlive = true, 1000, 1000))
92+
93+
val captured = EasyMock.newCapture[Select.Where]()
94+
expecting {
95+
factory.buildCluster(config).andReturn(cluster).once()
96+
cluster.connect().andReturn(session).once()
97+
keyspaceMetadata.getTable(tableName).andReturn(tableMetadata).once()
98+
metadata.getKeyspace(keyspaceName).andReturn(keyspaceMetadata).once()
99+
cluster.getMetadata.andReturn(metadata).once()
100+
session.getCluster.andReturn(cluster).once()
101+
session.prepare(EasyMock.capture(captured)).andReturn(selectPrepStatement).anyTimes()
102+
session.close().once()
103+
cluster.close().once()
104+
}
105+
106+
whenExecuting(factory, cluster, session, metadata, keyspaceMetadata, tableMetadata, selectPrepStatement) {
107+
val session = new CassandraSession(config, factory)
108+
session.ensureKeyspace(config.tracesKeyspace)
109+
val stmt = session.selectRawTracesPreparedStmt
110+
stmt shouldBe selectPrepStatement
111+
captured.getValue.getQueryString() shouldEqual "SELECT * FROM \"keyspace-1\".\"table-1\" WHERE id IN :id;"
112+
session.close()
113+
}
114+
}
72115
}
73116
}

pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,7 @@
6363
<properties>
6464
<project.jdk.version>1.8</project.jdk.version>
6565
<protobuf.version>3.4.0</protobuf.version>
66-
<haystack-commons.version>1.0.43</haystack-commons.version>
66+
<haystack-commons.version>1.0.46</haystack-commons.version>
6767

6868
<logback.version>1.2.3</logback.version>
6969
<slf4j-api.version>1.7.25</slf4j-api.version>

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -125,6 +125,12 @@ class TraceReader(traceStore: TraceStore,
125125
.getTraceCounts(request)
126126
}
127127

128+
def getRawTraces(request: RawTracesRequest): Future[RawTracesResult] = {
129+
traceStore
130+
.getRawTraces(request)
131+
.flatMap(traces => Future.successful(RawTracesResult.newBuilder().addAllTraces(traces.asJava).build()))
132+
}
133+
128134
private def buildTraceCallGraph(trace: Trace): TraceCallGraph = {
129135
val calls = trace.getChildSpansList
130136
.asScala

reader/src/main/scala/com/expedia/www/haystack/trace/reader/services/TraceService.scala

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,4 +126,10 @@ class TraceService(traceStore: TraceStore,
126126
traceReader.getTraceCounts(request)
127127
}
128128
}
129+
130+
override def getRawTraces(request: RawTracesRequest, responseObserver: StreamObserver[RawTracesResult]): Unit = {
131+
handleTraceCallGraphResponse.handle(request, responseObserver) {
132+
traceReader.getRawTraces(request)
133+
}
134+
}
129135
}

reader/src/main/scala/com/expedia/www/haystack/trace/reader/stores/CassandraEsTraceStore.scala

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,10 @@ class CassandraEsTraceStore(cassandraConfig: CassandraConfiguration,
144144
.map(result => mapSearchResultToTraceCount(request.getStartTime, request.getEndTime, result))
145145
}
146146

147+
override def getRawTraces(request: RawTracesRequest): Future[Seq[Trace]] = {
148+
cassandraReader.readRawTraces(request.getTraceIdList.asScala.toList)
149+
}
150+
147151
// convert all Futures to Try to make sure they all complete
148152
private def liftToTry[T](futures: Seq[Future[T]]): Seq[Future[Try[T]]] = futures.map { f =>
149153
f.map(Try(_)).recover { case t: Throwable => Failure(t) }

reader/src/main/scala/com/expedia/www/haystack/trace/reader/stores/TraceStore.scala

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,4 +26,5 @@ trait TraceStore extends AutoCloseable {
2626
def getFieldNames(): Future[Seq[String]]
2727
def getFieldValues(request: FieldValuesRequest): Future[Seq[String]]
2828
def getTraceCounts(request: TraceCountsRequest): Future[TraceCounts]
29+
def getRawTraces(request: RawTracesRequest): Future[Seq[Trace]]
2930
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,85 @@
1+
/*
2+
* Copyright 2017 Expedia, Inc.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
17+
package com.expedia.www.haystack.trace.reader.stores.readers.cassandra
18+
19+
import com.codahale.metrics.{Meter, Timer}
20+
import com.datastax.driver.core.exceptions.NoHostAvailableException
21+
import com.datastax.driver.core.{ResultSet, ResultSetFuture, Row}
22+
import com.expedia.open.tracing.api.Trace
23+
import com.expedia.www.haystack.commons.health.HealthController
24+
import com.expedia.www.haystack.trace.commons.clients.cassandra.CassandraTableSchema
25+
import com.expedia.www.haystack.trace.reader.exceptions.TraceNotFoundException
26+
import com.expedia.www.haystack.trace.reader.stores.readers.cassandra.CassandraReadRawTracesResultListener._
27+
import org.slf4j.{Logger, LoggerFactory}
28+
29+
import scala.collection.JavaConverters._
30+
import scala.collection.mutable
31+
import scala.concurrent.Promise
32+
import scala.util.{Failure, Success, Try}
33+
34+
object CassandraReadRawTracesResultListener {
35+
protected val LOGGER: Logger = LoggerFactory.getLogger(classOf[CassandraReadRawTracesResultListener])
36+
}
37+
38+
class CassandraReadRawTracesResultListener(asyncResult: ResultSetFuture,
39+
timer: Timer.Context,
40+
failure: Meter,
41+
promise: Promise[Seq[Trace]]) extends Runnable {
42+
override def run(): Unit = {
43+
timer.close()
44+
45+
Try(asyncResult.get)
46+
.flatMap(tryGetTraceRows)
47+
.flatMap(tryDeserialize)
48+
match {
49+
case Success(traces) =>
50+
promise.success(traces)
51+
case Failure(ex) =>
52+
if (fatalError(ex)) {
53+
LOGGER.error("Fatal error in reading from cassandra, tearing down the app", ex)
54+
HealthController.setUnhealthy()
55+
} else {
56+
LOGGER.error("Failed in reading the record from cassandra", ex)
57+
}
58+
failure.mark()
59+
promise.failure(ex)
60+
}
61+
}
62+
63+
private def fatalError(ex: Throwable): Boolean = {
64+
if (ex.isInstanceOf[NoHostAvailableException]) true else ex.getCause != null && fatalError(ex.getCause)
65+
}
66+
67+
private def tryGetTraceRows(resultSet: ResultSet): Try[Seq[Row]] = {
68+
val rows = resultSet.all().asScala
69+
if (rows.isEmpty) Failure(new TraceNotFoundException) else Success(rows)
70+
}
71+
72+
private def tryDeserialize(rows: Seq[Row]): Try[Seq[Trace]] = {
73+
val traceBuilderMap = new mutable.HashMap[String, Trace.Builder]()
74+
var deserFailed: Failure[Seq[Trace]] = null
75+
76+
rows.foreach(row => {
77+
CassandraTableSchema.extractSpanBufferFromRow(row) match {
78+
case Success(sBuffer) =>
79+
traceBuilderMap.getOrElseUpdate(sBuffer.getTraceId, Trace.newBuilder().setTraceId(sBuffer.getTraceId)).addAllChildSpans(sBuffer.getChildSpansList)
80+
case Failure(cause) => deserFailed = Failure(cause)
81+
}
82+
})
83+
if (deserFailed == null) Success(traceBuilderMap.values.map(_.build).toSeq) else deserFailed
84+
}
85+
}

reader/src/main/scala/com/expedia/www/haystack/trace/reader/stores/readers/cassandra/CassandraTraceReader.scala

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,5 +49,23 @@ class CassandraTraceReader(cassandra: CassandraSession, config: CassandraConfigu
4949
}
5050
}
5151

52+
def readRawTraces(traceIds: List[String]): Future[Seq[Trace]] = {
53+
val timer = readTimer.time()
54+
val promise = Promise[Seq[Trace]]
55+
56+
try {
57+
val statement = cassandra.newSelectRawTracesBoundStatement(traceIds)
58+
val asyncResult = cassandra.executeAsync(statement)
59+
asyncResult.addListener(new CassandraReadRawTracesResultListener(asyncResult, timer, readFailures, promise), dispatcher)
60+
promise.future
61+
} catch {
62+
case ex: Exception =>
63+
readFailures.mark()
64+
timer.stop()
65+
LOGGER.error("Failed to read raw traces with exception", ex)
66+
Future.failed(ex)
67+
}
68+
}
69+
5270
override def close(): Unit = ()
5371
}

reader/src/test/scala/com/expedia/www/haystack/trace/reader/integration/TraceServiceIntegrationTestSpec.scala

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -393,4 +393,29 @@ class TraceServiceIntegrationTestSpec extends BaseIntegrationTestSpec {
393393
traceCounts.getTraceCountList.asScala.foreach(_.getCount shouldBe 1)
394394
}
395395
}
396+
397+
describe("TraceReader.getRawTraces") {
398+
it("should get raw traces for given traceIds from cassandra") {
399+
Given("traces in cassandra")
400+
val traceId1 = UUID.randomUUID().toString
401+
val spanId1 = UUID.randomUUID().toString
402+
val spanId2 = UUID.randomUUID().toString
403+
putTraceInCassandra(traceId1, spanId1, "svc1", "oper1")
404+
putTraceInCassandra(traceId1, spanId2, "svc2", "oper2")
405+
406+
val traceId2 = UUID.randomUUID().toString
407+
val spanId3 = UUID.randomUUID().toString
408+
putTraceInCassandra(traceId2, spanId3, "svc1", "oper1")
409+
410+
When("getRawTraces is invoked")
411+
val tracesResult = client.getRawTraces(RawTracesRequest.newBuilder().addAllTraceId(Seq(traceId1, traceId2).asJava).build())
412+
413+
Then("should return the traces")
414+
val traceIdSpansMap: Map[String, Set[String]] = tracesResult.getTracesList.asScala
415+
.map(trace => (trace.getTraceId -> trace.getChildSpansList().asScala.map(_.getSpanId).toSet)).toMap
416+
417+
traceIdSpansMap(traceId1) shouldEqual (Set(spanId1, spanId2))
418+
traceIdSpansMap(traceId2) shouldEqual (Set(spanId3))
419+
}
420+
}
396421
}

0 commit comments

Comments
 (0)