@@ -21,7 +21,7 @@ import com.expedia.www.haystack.trace.commons.clients.cassandra.{CassandraCluste
2121import com .expedia .www .haystack .trace .commons .clients .es .document .TraceIndexDoc
2222import com .expedia .www .haystack .trace .commons .config .entities .{CassandraConfiguration , WhitelistIndexFieldConfiguration }
2323import com .expedia .www .haystack .trace .reader .config .entities .{ElasticSearchConfiguration , ServiceMetadataReadConfiguration }
24- import com .expedia .www .haystack .trace .reader .metrics .{ AppMetricNames , MetricsSupport }
24+ import com .expedia .www .haystack .trace .reader .metrics .MetricsSupport
2525import com .expedia .www .haystack .trace .reader .stores .readers .ServiceMetadataReader
2626import com .expedia .www .haystack .trace .reader .stores .readers .cassandra .CassandraTraceReader
2727import com .expedia .www .haystack .trace .reader .stores .readers .es .ElasticSearchReader
@@ -32,7 +32,6 @@ import org.slf4j.LoggerFactory
3232
3333import scala .collection .JavaConverters ._
3434import scala .concurrent .{ExecutionContextExecutor , Future }
35- import scala .util .{Failure , Success , Try }
3635
3736class CassandraEsTraceStore (cassandraConfig : CassandraConfiguration ,
3837 serviceMetadataConfig : ServiceMetadataReadConfiguration ,
@@ -41,7 +40,6 @@ class CassandraEsTraceStore(cassandraConfig: CassandraConfiguration,
4140 extends TraceStore with MetricsSupport with ResponseParser {
4241
4342 private val LOGGER = LoggerFactory .getLogger(classOf [ElasticSearchReader ])
44- private val traceRejected = metricRegistry.meter(AppMetricNames .SEARCH_TRACE_REJECTED )
4543
4644 private val cassandraSession = new CassandraSession (cassandraConfig, new CassandraClusterFactory )
4745 private val cassandraReader : CassandraTraceReader = new CassandraTraceReader (cassandraSession, cassandraConfig)
@@ -79,36 +77,21 @@ class CassandraEsTraceStore(cassandraConfig: CassandraConfiguration,
7977 // go through each hit and fetch trace for parsed traceId
8078 val sourceList = result.getSourceAsStringList
8179 if (sourceList != null && sourceList.size() > 0 ) {
82- val traceFutures = sourceList
80+ val traceIds = sourceList
8381 .asScala
8482 .map(source => extractTraceIdFromSource(source))
8583 .filter(! _.isEmpty)
86- .toSet[String ]
87- .toSeq
88- .map(id => getTrace(id))
89-
90- // wait for all Futures to complete and then map them to Traces
91- Future
92- .sequence(liftToTry(traceFutures))
93- .map(_.flatMap(retrieveTriedTrace))
84+ .toSet[String ] // de-dup traceIds
85+ .toList
86+
87+ cassandraReader.readRawTraces(traceIds)
9488 } else {
9589 Future .successful(Nil )
9690 }
9791 }
9892
9993 override def getTrace (traceId : String ): Future [Trace ] = cassandraReader.readTrace(traceId)
10094
101- private def retrieveTriedTrace (mayBeTrace : Try [Trace ]): Option [Trace ] = {
102- mayBeTrace match {
103- case Success (trace) =>
104- Some (trace)
105- case Failure (ex) =>
106- LOGGER .warn(" traceId not found in cassandra, rejected searched trace" , ex)
107- traceRejected.mark()
108- None
109- }
110- }
111-
11295 override def getFieldNames (): Future [Seq [String ]] = {
11396 Future .successful(indexConfig.whitelistIndexFields.map(_.name).distinct.sorted)
11497 }
@@ -148,11 +131,6 @@ class CassandraEsTraceStore(cassandraConfig: CassandraConfiguration,
148131 cassandraReader.readRawTraces(request.getTraceIdList.asScala.toList)
149132 }
150133
151- // convert all Futures to Try to make sure they all complete
152- private def liftToTry [T ](futures : Seq [Future [T ]]): Seq [Future [Try [T ]]] = futures.map { f =>
153- f.map(Try (_)).recover { case t : Throwable => Failure (t) }
154- }
155-
156134 override def close (): Unit = {
157135 cassandraReader.close()
158136 esReader.close()
0 commit comments