Skip to content
This repository has been archived by the owner on Jan 20, 2022. It is now read-only.

Fixing tormentaSpout.registerMetric call in scheduleSpout #521

Merged
merged 1 commit into from
Jun 23, 2014
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -216,9 +216,8 @@ abstract class Storm(options: Map[String, Options], transformConfig: Summingbird
}

val metrics = getOrElse(stormDag, node, DEFAULT_SPOUT_STORM_METRICS)
tormentaSpout.registerMetrics(metrics.toSpoutMetrics)

val stormSpout = tormentaSpout.getSpout
val stormSpout = tormentaSpout.registerMetrics(metrics.toSpoutMetrics).getSpout
val parallelism = getOrElse(stormDag, node, parOpt.getOrElse(DEFAULT_SPOUT_PARALLELISM)).parHint
topologyBuilder.setSpout(nodeName, stormSpout, parallelism)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,10 @@ limitations under the License.

package com.twitter.summingbird.storm.option

import com.twitter.summingbird.storm.StormMetric
import backtype.storm.metric.api.IMetric
import com.twitter.summingbird.storm.StormMetric
import com.twitter.tormenta.spout.Metric
import java.io.Serializable

/**
* Options used by the flatMapping stage of a storm topology.
Expand Down Expand Up @@ -56,7 +57,7 @@ object SpoutStormMetrics {
def unapply(metrics: SpoutStormMetrics) = Some(metrics.metrics)
}

class SpoutStormMetrics(val metrics: () => TraversableOnce[StormMetric[IMetric]]) {
class SpoutStormMetrics(val metrics: () => TraversableOnce[StormMetric[IMetric]]) extends Serializable {
def toSpoutMetrics: () => TraversableOnce[Metric[IMetric]] =
{ () => metrics().map { x: StormMetric[IMetric] => Metric(x.name, x.metric, x.interval.inSeconds) } }
}
Expand Down