亚洲在线久爱草,狠狠天天香蕉网,天天搞日日干久草,伊人亚洲日本欧美

為了賬號安全,請及時綁定郵箱和手機立即綁定

Storm框架:如何消費RabbitMq消息(代碼案例)

標簽:
Java

1、定义拓扑topology

public class MessageTopology {    public static void main(String[] args) throws Exception {        //组装topology
        TopologyBuilder topologyBuilder = new TopologyBuilder();
        topologyBuilder.setSpout("RabbitmqSpout", new RabbitmqSpout());
        topologyBuilder.setBolt("FilterBolt", new FilterBolt()).shuffleGrouping("RabbitmqSpout");

        Config conf = new Config ();        try {            if (args.length > 0) {
                StormSubmitter.submitTopology(args[0], conf, topologyBuilder.createTopology());
            } else {
                LocalCluster localCluster = new LocalCluster();
                localCluster.submitTopology("messageTopology", conf, topologyBuilder.createTopology());
            }
        } catch (AlreadyAliveException e) {
            e.printStackTrace();
        }
    }
}

2、定义数据源RabbitmqSpout

RabbitmqSpout继承自org.apache.storm.topology.IRichSpout接口,实现对应的方法:open(),close(),activate(),deactivate(),nextTuple(),ack(),fail()。

unconfirmedMap对象存储了MQ所有发射出去等待确认的消息唯一标识deliveryTag,当storm系统回调ack、fail方法后进行MQ消息的成功确认或失败重回队列操作(Storm系统回调方法会在bolt操作中主动调用ack、fail方法时触发)。

public class RabbitmqSpout implements IRichSpout {    private final Logger LOGGER = LoggerFactory.getLogger(RabbitmqSpout.class);    private Map map;    private TopologyContext topologyContext;    private SpoutOutputCollector spoutOutputCollector;    private Connection connection;    private Channel channel;    private static final String QUEUE_NAME = "message_queue";    private final Map<String, Long> unconfirmedMap = Collections.synchronizedMap(new HashMap<String, Long>());    //连接mq服务
    private void connect() throws IOException, TimeoutException {
        ConnectionFactory factory = new ConnectionFactory();
        factory.setHost("127.0.0.1");
        factory.setUsername("admin");
        factory.setPassword("admin");
        factory.setVirtualHost("/");

        connection = factory.newConnection();
        channel = connection.createChannel();
        channel.queueDeclare(QUEUE_NAME, true, false, false, null);
    }    @Override
    public void open(Map map, TopologyContext topologyContext, SpoutOutputCollector spoutOutputCollector) {        this.map = map;        this.topologyContext = topologyContext;        this.spoutOutputCollector = spoutOutputCollector;        try {            this.connect();
        } catch (IOException e) {
            e.printStackTrace();
        } catch (TimeoutException e) {
            e.printStackTrace();
        }
    }    @Override
    public void close() {        try {
            channel.close();
            connection.close();
        } catch (IOException e) {
            e.printStackTrace();
        } catch (TimeoutException e) {
            e.printStackTrace();
        }
    }    @Override
    public void nextTuple() {        try {
            GetResponse response = channel.basicGet(QUEUE_NAME, false);            if (response == null) {
                Utils.sleep(3000);
            } else {
                AMQP.BasicProperties props = response.getProps();
                String messageId = UUID.randomUUID().toString();
                Long deliveryTag = response.getEnvelope().getDeliveryTag();
                String body = new String(response.getBody());

                unconfirmedMap.put(messageId, deliveryTag);
                LOGGER.info("RabbitmqSpout: {}, {}, {}, {}", body, messageId, deliveryTag, props);                this.spoutOutputCollector.emit(new Values(body), messageId);
            }
        } catch (IOException e) {
            e.printStackTrace();
        }
    }    @Override
    public void ack(Object o) {
        String messageId = o.toString();
        Long deliveryTag = unconfirmedMap.get(messageId);
        LOGGER.info("ack: {}, {}, {}\n\n", messageId, deliveryTag, unconfirmedMap.size());        try {
            unconfirmedMap.remove(messageId);
            channel.basicAck(deliveryTag, false);
        } catch (IOException e) {
            e.printStackTrace();
        }
    }    @Override
    public void fail(Object o) {
        String messageId = o.toString();
        Long deliveryTag = unconfirmedMap.get(messageId);
        LOGGER.info("fail: {}, {}, {}\n\n", messageId, deliveryTag, unconfirmedMap.size());        try {
            unconfirmedMap.remove(messageId);
            channel.basicNack(deliveryTag, false, true);
        } catch (IOException e) {
            e.printStackTrace();
        }
    }    @Override
    public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
        outputFieldsDeclarer.declare(new Fields("body"));
    }    @Override
    public Map<String, Object> getComponentConfiguration() {        return null;
    }    
    @Override
    public void activate() {

    }    @Override
    public void deactivate() {

    }
}

3、定义数据流处理FilterBolt

public class FilterBolt implements IRichBolt {    private final Logger LOGGER = LoggerFactory.getLogger(FilterBolt.class);    private Map map;    private TopologyContext topologyContext;    private OutputCollector outputCollector;    @Override
    public void prepare(Map map, TopologyContext topologyContext, OutputCollector outputCollector) {        this.map = map;        this.topologyContext = topologyContext;        this.outputCollector = outputCollector;
    }    @Override
    public void execute(Tuple tuple) {
        String value = tuple.getStringByField("body");

        LOGGER.info("FilterBolt:{}", value);
        outputCollector.ack(tuple);
    }    
    @Override
    public void declareOutputFields(OutputFieldsDeclarer outputFieldsDeclarer) {
        outputFieldsDeclarer.declare(new Fields("body"));
    }    @Override
    public Map<String, Object> getComponentConfiguration() {        return null;
    }    
    @Override
    public void cleanup() {

    }
}

Hey, show me the code!

原文出处:https://www.cnblogs.com/gouyg/p/java_storm_rabbitmq_example.html  

點擊查看更多內容
TA 點贊

若覺得本文不錯,就分享一下吧!

評論

作者其他優質文章

正在加載中
  • 推薦
  • 1
  • 收藏
  • 共同學習,寫下你的評論
感謝您的支持,我會繼續努力的~
掃碼打賞,你說多少就多少
贊賞金額會直接到老師賬戶
支付方式
打開微信掃一掃,即可進行掃碼打賞哦
今天注冊有機會得

100積分直接送

付費專欄免費學

大額優惠券免費領

立即參與 放棄機會
微信客服

購課補貼
聯系客服咨詢優惠詳情

幫助反饋 APP下載

慕課網APP
您的移動學習伙伴

公眾號

掃描二維碼
關注慕課網微信公眾號

舉報

0/150
提交
取消