基礎storm程式示例
Storm的流處理主要就是通過Spout和Bolt節點進行處理,可以繼承這些類寫自己的邏輯
public class FlinkStormDemo {
public static void main(String[] args) {
//1.創建執行環境
LocalCluster stormCluster = new LocalCluster();
TopologyBuilder builder = new TopologyBuilder();
//2.創建一個初始的資料源
builder.setSpout("word", new WordSpout());
//3.對資料源進行第一次加工
builder.setBolt("word-1",new WordBolt1(), 1).shuffleGrouping("word");
//4.對資料源進行第二次加工
builder.setBolt("word-2",new WordBolt2(), 1).shuffleGrouping("word-1");
//5.配置一些引數
Config config = new Config();
config.setDebug(true);
//6.提交storm任務并處理
stormCluster.submitTopology("storm-task", config, builder.createTopology());
}
static class WordSpout extends BaseRichSpout {
private SpoutOutputCollector spoutOutputCollector;
@Override
public void open(Map map, TopologyContext topologyContext, SpoutOutputCollector spoutOutputCollector) {
this.spoutOutputCollector = spoutOutputCollector;
}
@Override
public void nextTuple() {
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("資料初始化中....");
String initData = "abc";
spoutOutputCollector.emit(new Values(initData));
}
@Override
public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
outputFieldsDeclarer.declare(new Fields("word"));
}
}
static class WordBolt1 extends BaseRichBolt {
private OutputCollector collector;
@Override
public void prepare(Map map, TopologyContext topologyContext, OutputCollector outputCollector) {
this.collector = outputCollector;
}
@Override
public void execute(Tuple tuple) {
System.out.println("資料第1次處理中....");
//給上次獲取的單詞拼接上def
collector.emit(tuple, new Values(tuple.getString(0) + "def"));
collector.ack(tuple);
}
@Override
public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
outputFieldsDeclarer.declare(new Fields("word"));
}
}
static class WordBolt2 extends BaseRichBolt {
private OutputCollector collector;
@Override
public void prepare(Map map, TopologyContext topologyContext, OutputCollector outputCollector) {
this.collector = outputCollector;
}
@Override
public void execute(Tuple tuple) {
System.out.println("資料第2次處理中....");
//輸出處理結果
System.out.println("處理結果:" + tuple.getString(0));
collector.ack(tuple);
}
@Override
public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
}
}
}
執行結果:

利用flink-storm程式實作類似功能
需要更改flink相關依賴的版本到1.7.0,主要依賴了flink-storm的jar包
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>1.7.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_2.11</artifactId>
<version>1.7.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients_2.11</artifactId>
<version>1.7.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-storm_2.11</artifactId>
<version>1.7.0</version>
</dependency>
只需要修改兩處:
LocalCluster替換為FlinkLocalCluster,處理的任務從TopologyBuilder
.createTopology替換為FlinkTopology.createTopology(TopologyBuilder)

執行結果:

利用flink程式實作類似功能
利用Kafka發送初始訊息“測驗資料”,
public class FlinkProducer {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
Properties properties = new Properties();
properties.put("bootstrap.servers", "10.225.173.107:9092,10.225.173.108:9092,10.225.173.109:9092");
FlinkKafkaProducer<String> flinkKafkaProducer = new FlinkKafkaProducer<>("flink", new SimpleStringSchema(), properties);
DataStreamSource<String> source = env.fromElements("測驗資料");
source.addSink(flinkKafkaProducer);
env.execute();
}
}
接收Kafka的初始訊息“測驗資料”并加工處理,
public class FlinkConsumer {
public static void main(String[] args) throws Exception {
//1.創建執行環境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
//2.配置并創建初始資料源
Properties properties = new Properties();
properties.put("bootstrap.servers", "10.225.173.107:9092,10.225.173.108:9092,10.225.173.109:9092");
FlinkKafkaConsumer<String> flinkKafkaConsumer = new FlinkKafkaConsumer<>("flink", new SimpleStringSchema(), properties);
DataStreamSource<String> source = env.addSource(flinkKafkaConsumer);
//3.對資料源進行連續處理
source.process(new FlinkBolt1()).
process(new FlinkBolt2()).
process(new FlinkBolt3());
//4.執行flink程式
env.execute();
}
static class FlinkBolt1 extends ProcessFunction<String, Object> {
@Override
public void open(Configuration parameters) {
//開始第1次處理....
}
@Override
public void processElement(String s, ProcessFunction<String, Object>.Context context, Collector<Object> collector) throws Exception {
System.out.println("第1次處理前的值是:" + s);
s += "abc";
System.out.println("第1次處理后的值是:" + s);
collector.collect(s);
}
@Override
public void close() throws Exception {
//結束第1次處理
}
}
static class FlinkBolt2 extends ProcessFunction<Object, Object> {
@Override
public void open(Configuration parameters) {
//開始第2次處理....
}
@Override
public void processElement(Object s, ProcessFunction<Object, Object>.Context context, Collector<Object> collector) throws Exception {
s = s.toString();
System.out.println("第2次處理前的值是:" + s);
s += "def";
System.out.println("第2次處理后的值是:" + s);
collector.collect(s);
}
@Override
public void close() throws Exception {
//結束第2次處理....
}
}
static class FlinkBolt3 extends ProcessFunction<Object, Object> {
@Override
public void open(Configuration parameters) {
//開始第3次處理....
}
@Override
public void processElement(Object s, ProcessFunction<Object, Object>.Context context, Collector<Object> collector) throws Exception {
s = s.toString();
System.out.println("第3次處理前的值是:" + s);
s += "ghi";
System.out.println("第3次處理后的值是:" + s);
collector.collect(s);
}
@Override
public void close() throws Exception {
//結束第3次處理....
}
}
}
處理結果:

轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/436405.html
標籤:其他
上一篇:Hadoop13:【案例】MapReduce任務日志查看
下一篇:HBase 過濾器
