java 实现 apollo作为消息代理服务器
Apache apollo(阿波罗)是消息代理服务器,实在activeMQ基础上来的,可以支持MQTT,WebSokcet等多种协议。
MQTT(消息队列遥测传输)是ISO 标准(ISO/IEC PRF 20922)下基于发布/订阅范式的消息协议。它工作在TCP/IP协议族上,是为硬件性能低下的远程设备以及网络状况糟糕的情况下而设计的发布/订阅型消息协议,为此,它需要一个消息中间件,现在兴起的智慧工厂需要将仪表,设备等数据发送到云端,MQTT协议无疑是有自己独特的优点。
下面介绍一下,如何通过阿波罗消息代理服务器实现MQTT协议的发布订阅,比如,在工厂里面有多个仪表需要将仪表数据上传的服务器来进行实现。
简单描述一下,各个仪表通过仪表自身协议上传到网关,由网关通过MQTT协议将数据上传到服务器中,由服务器端的程序进行解析,存库。MQTT为发布订阅协议,关于发布订阅,举一个简单的例子。公司老板要公布放假消息和加班消息,公司老板不可能一个个的去告诉员工,比如他会让相关部门人员在公告板上发布消息通知。员工1在公告板上发布了放假通知,员工2在公告板上发布了加班通知,然后员工3订阅了放假通知的主题,员工4订阅了加班通知,那么员工3会收到放假通知,员工4会收到加班通知。这样对应过去,网关就是作为消息发布者,阿波罗代理服务器就相当于公告板,我们需要去订阅网关转发的消息主题,这样我们就能上收到网关转发的数据。
阿波罗的安装:
一、下载apollo,因为官网上经过改动不是很好找资源,所以我已经上传到了****。
二、解压到电脑上,我用的是win10,所以下载对应版本,解压后目录结构。
三、以管理员的方式打开控制台,进入到安装apollo的bin目录下,创建broker, apollo create myapollo+创建broker的文件路径+broker
四、运行broker 在控制台上输入apollo-broker run
五:输入http://127.0.0.1:61680/console/index.html,然后输入默认的用户名admin和密码password即可进入。
六、基本配置在创建的broker的etc的apollo.xml中进行配置,比如需要将tcp端口61613映射到外网,可以在此修改。
基于JAVA通过阿波罗实现MQTT协议
1:jar包org.eclipse.paho.client.mqttv3-1.1.0.jar
2:订阅段代码
public class ReportMqtt implements MqttCallback { public String HOST="tcp://127.0.0.1:61613" ; public String TOPIC="123"; private String name="admin"; private String passWord ="password"; private MqttClient client; private MqttConnectOptions options; private MqttMessage message; String clientid= UUID.randomUUID().toString(); private static Logger logger=Logger.getLogger(ReportMqtt.class); //订阅 public void subscribe() { try { // host为主机名,clientid即连接MQTT的客户端ID,一般以唯一标识符表示,MemoryPersistence设置clientid的保存形式,默认为以内存保存 client = new MqttClient(HOST, clientid, new MemoryPersistence()); // MQTT的连接设置 options = new MqttConnectOptions(); // 设置是否清空session,这里如果设置为false表示服务器会保留客户端的连接记录,这里设置为true表示每次连接到服务器都以新的身份连接 options.setCleanSession(true); // 设置连接的用户名 options.setUserName(name); // 设置连接的密码 options.setPassword(passWord.toCharArray()); // 设置超时时间 单位为秒 options.setConnectionTimeout(10); // 设置会话心跳时间 单位为秒 服务器会每隔1.5*20秒的时间向客户端发送个消息判断客户端是否在线,但这个方法并没有重连的机制 options.setKeepAliveInterval(3600); // 设置回调 client.setCallback(this); client.connect(options); //订阅消息 int[] Qos = {1}; String[] topic1 = {TOPIC}; client.subscribe(topic1, Qos); System.out.println("订阅端订取消息成功"); } catch (Exception e) { System.out.println("ReportMqtt客户端连接异常,异常信息:"+e); } } @Override public void connectionLost(Throwable throwable) { try { System.out.println("程序出现异常,DReportMqtt断线!正在重新连接...:"); client.close(); //重新订阅 this.subscribe(); System.out.println("ReportMqtt重新连接成功"); }catch (MqttException e){ System.out.println(e.getMessage()); } } @Override public void messageArrived(String topic, MqttMessage mqttMessage) throws Exception { System.out.println("接收消息主题:"+ topic); System.out.println("接收消息Qos:"+ mqttMessage.getQos()); System.out.println("接收消息内容 :"+ new String(mqttMessage.getPayload())); } @Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println("消息发送成功"); } }
u
3:发布端代码
ord="password";
public class PlatformMqtt implements MqttCallback { private String HOST="tcp://127.0.0.1:61613"; private String name="admin"; private String passWord="password"; private MqttClient client; private MqttMessage message; String clientid= UUID.randomUUID().toString(); public void send() { try { client = new MqttClient(HOST, clientid, new MemoryPersistence()); } catch (MqttException e) { e.printStackTrace(); } MqttConnectOptions options = new MqttConnectOptions(); options.setCleanSession(false); options.setUserName(name); options.setPassword(passWord.toCharArray()); // 设置超时时间 options.setConnectionTimeout(10); // 设置会话心跳时间 options.setKeepAliveInterval(3600); try { client.setCallback((MqttCallback) this); client.connect(options); System.out.println("客户端连接阿波罗成功"); } catch (Exception e) { System.out.println("platform-Mqtt客户端连接异常,异常信息:"+e); } } @RequestMapping("/request") @ResponseBody public void publish() throws MqttException { System.out.println("客户端开始发布消息"); String body="能发消息不"; String topic ="123"; message = new MqttMessage(); message.setQos(0); message.setRetained(false); message.setPayload(body.getBytes()); client.publish(topic,message); System.out.println("客户端发部消息成功"); } @Override public void connectionLost(Throwable throwable) { System.out.println("程序出现异常,正在重新连接...:"); try { client.close(); send(); } catch (MqttException e) { System.out.println(e.getMessage()); } System.out.println("platform-Mqtt重新连接成功"); } @Override public void messageArrived(String s, MqttMessage mqttMessage) { } @Override public void deliveryComplete(IMqttDeliveryToken iMqttDeliveryToken) { System.out.println("消息发送成功"); } }
private MqttCliet client;