/
ChannelViewTimeBolt.java
52 lines (44 loc) · 1.4 KB
/
ChannelViewTimeBolt.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
package org.nerdronix.stormchronicle;
import java.util.HashMap;
import java.util.Map;
import backtype.storm.task.OutputCollector;
import backtype.storm.task.TopologyContext;
import backtype.storm.topology.OutputFieldsDeclarer;
import backtype.storm.topology.base.BaseRichBolt;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Tuple;
import backtype.storm.tuple.Values;
/**
* This aggregator class will calculate total amount of time all users spent
* viewing a channel
*
* @author Abraham Menacherry
*
*/
public class ChannelViewTimeBolt extends BaseRichBolt {
private static final long serialVersionUID = 1L;
private OutputCollector collector;
private Map<Integer, Integer> channelViewTimeMap = new HashMap<>();
@Override
public void prepare(@SuppressWarnings("rawtypes") Map stormConf,
TopologyContext context, OutputCollector collector) {
this.collector = collector;
}
@Override
public void execute(Tuple input) {
Integer channelId = input.getInteger(1);
Integer viewTime = channelViewTimeMap.get(channelId);
if (null == viewTime) {
viewTime = input.getInteger(2);
} else {
viewTime += input.getInteger(2);
}
channelViewTimeMap.put(channelId, viewTime);
collector.emit(new Values(channelId, viewTime));
collector.ack(input);
}
@Override
public void declareOutputFields(OutputFieldsDeclarer declarer) {
declarer.declare(new Fields("channelId", "viewTime"));
}
}