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
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
| // Storm拓扑
import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.StormSubmitter;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichSpout;
import org.apache.storm.topology.base.BaseBasicBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;
import org.apache.storm.tuple.Tuple;
public class WordCountTopology {
public static void main(String[] args) throws Exception {
TopologyBuilder builder = new TopologyBuilder();
// Spout: 数据源
builder.setSpout("sentence-spout", new SentenceSpout(), 2);
// Bolt: 分词
builder.setBolt("split-bolt", new SplitBolt(), 4)
.setNumTasks(8)
.shuffleGrouping("sentence-spout");
// Bolt: 计数
builder.setBolt("count-bolt", new CountBolt(), 4)
.fieldsGrouping("split-bolt", new Fields("word"));
// Bolt: 报告
builder.setBolt("report-bolt", new ReportBolt(), 1)
.globalGrouping("count-bolt");
Config config = new Config();
config.setDebug(true);
if (args != null && args.length > 0) {
// 生产集群
config.setNumWorkers(4);
StormSubmitter.submitTopologyWithProgressBar(
args[0],
config,
builder.createTopology()
);
} else {
// 本地测试
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("word-count", config,
builder.createTopology());
Thread.sleep(60000);
cluster.shutdown();
}
}
// Sentence Spout
public static class SentenceSpout extends BaseRichSpout {
private SpoutOutputCollector collector;
private String[] sentences = {
"the quick brown fox",
"jumps over the lazy dog",
"hello world from storm"
};
@Override
public void open(
Map conf,
TopologyContext context,
SpoutOutputCollector collector
) {
this.collector = collector;
}
@Override
public void nextTuple() {
for (String sentence : sentences) {
collector.emit(new Values(sentence));
}
Utils.sleep(100);
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
declarer.declare(new Fields("sentence"));
}
}
// Split Bolt
public static class SplitBolt extends BaseBasicBolt {
@Override
public void execute(Tuple tuple, BasicOutputCollector collector) {
String sentence = tuple.getStringByField("sentence");
String[] words = sentence.split(" ");
for (String word : words) {
collector.emit(new Values(word));
}
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
declarer.declare(new Fields("word"));
}
}
// Count Bolt
public static class CountBolt extends BaseBasicBolt {
private Map<String, Integer> counts = new HashMap<>();
@Override
public void execute(Tuple tuple, BasicOutputCollector collector) {
String word = tuple.getStringByField("word");
Integer count = counts.getOrDefault(word, 0) + 1;
counts.put(word, count);
collector.emit(new Values(word, count));
}
@Override
public void declareOutputFields(
OutputFieldsDeclarer declarer
) {
declarer.declare(new Fields("word", "count"));
}
}
}
|