Created by vladimir on 17 14 import backtype storm Config import backt

 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
/**
* Created by vladimir on 4/17/14.
*/
import backtype.storm.Config;
import backtype.storm.LocalCluster;
import backtype.storm.tuple.Fields;
import backtype.storm.LocalDRPC;
import backtype.storm.generated.StormTopology;
import backtype.storm.tuple.Values;
import storm.kafka.KafkaConfig;
import storm.kafka.trident.TransactionalTridentKafkaSpout;
import storm.kafka.trident.TridentKafkaConfig;
import storm.trident.TridentTopology;
import storm.trident.operation.BaseFunction;
import storm.trident.operation.TridentCollector;
import storm.trident.operation.builtin.Count;
import storm.trident.tuple.TridentTuple;
import storm.trident.spout.IBatchSpout;
import java.io.File;
import java.util.Arrays;
import java.util.Map;
import java.util.ArrayList;
import backtype.storm.task.TopologyContext;
import java.io.FileInputStream;
import java.io.BufferedReader;
import java.io.InputStreamReader;
import storm.trident.*;
import storm.kafka.*;
public class MyReader {
public static class DebugBytes extends BaseFunction{
Integer count=0;
@Override
public void execute(TridentTuple tuple, TridentCollector collector){
count++;
if (count % 2 == 0){
throw new RuntimeException("Testing");
}
String text = new String (tuple.getBinary(0));
System.out.println(text);
collector.emit(new Values(text));
}
}
public static StormTopology makeTopology(LocalDRPC drpc){
TridentTopology topology = new TridentTopology();
TridentKafkaConfig spoutConfig = new TridentKafkaConfig(
KafkaConfig.StaticHosts.fromHostString(
Arrays.asList(new String[]{
"localhost"}), 2), "test");
topology.newStream("kafka", new TransactionalTridentKafkaSpout(spoutConfig)) .each(new Fields("bytes"), new DebugBytes(), new
Fields("text"));
return topology.build();
}
public static void main(String[] args) throws Exception{
/*Config conf = new Config();
conf.setDebug(true);
conf.setMaxTaskParallelism(3);
conf.setMaxSpoutPending(20);
conf.put(Config.STORM_ZOOKEEPER_SERVERS, Arrays.asList(new
String[]{"127.0.0.1"}));
conf.put(Config.STORM_ZOOKEEPER_PORT, 2181);
conf.put(Config.STORM_ZOOKEEPER_ROOT, "/storm");
LocalDRPC drpc = new LocalDRPC();
LocalCluster cluster = new LocalCluster();*/
Config conf = new Config();
conf.setMaxSpoutPending(20);
LocalDRPC drpc = new LocalDRPC();
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("transactional-topology", conf, makeTopology(drpc));
Thread.sleep(10000);
}
}