RabbitMQ学习之ConntectionFactory与Conntection的认知
在发送和接收消息重要的类有:ConnectionFactory, Connection,Channel和 QueueingConsumer。
ConntectionFactory类是方便创建与AMQP代理相关联的Connection;下面来看看ConntectionFactory是如何创建一个Contention.首先通过new ConnectionFactory()创建一个ConnectionFactory;并设置此连接工厂的主机设置为broke IP。通过ConnectionFactory的newConnection()方法 创建一个Connection; newConnection方法通过得到当前连接的地址及端口号来获得一个Address,通过createFrameHandler的方法 来得到FrameHandler;再通过new AMQConnection(this, frameHandler)来得到Connection并启动。
代码清单1创建Connection的源码(ConnectionFactory.Java中)
- protected FrameHandler createFrameHandler(Address addr)
- throws IOException {
- String hostName = addr.getHost();//获取主机IP
- int portNumber = portOrDefault(addr.getPort());//获取端口
- Socket socket = null;
- try {
- socket = factory.createSocket();
- configureSocket(socket);
- socket.connect(new InetSocketAddress(hostName, portNumber),
- connectionTimeout);
- return createFrameHandler(socket);
- } catch (IOException ioe) {
- quietTrySocketClose(socket);
- throw ioe;
- }
- }
- private static void quietTrySocketClose(Socket socket) {
- if (socket != null)
- try { socket.close(); } catch (Exception _) {/*ignore exceptions*/}
- }
- protected FrameHandler createFrameHandler(Socket sock)
- throws IOException
- {
- return new SocketFrameHandler(sock);
- }
- /**
- * Provides a hook to insert custom configuration of the sockets
- * used to connect to an AMQP server before they connect.
- *
- * The default behaviour of this method is to disable Nagle's
- * algorithm to get more consistently low latency. However it
- * may be overridden freely and there is no requirement to retain
- * this behaviour.
- *
- * @param socket The socket that is to be used for the Connection
- */
- protected void configureSocket(Socket socket) throws IOException{
- // disable Nagle's algorithm, for more consistently low latency
- socket.setTcpNoDelay(true);
- }
- /**
- * Create a new broker connection
- * @param addrs an array of known broker addresses (hostname/port pairs) to try in order
- * @return an interface to the connection
- * @throws IOException if it encounters a problem
- */
- public Connection newConnection(Address[] addrs) throws IOException {
- return newConnection(null, addrs);
- }
- /**
- * Create a new broker connection
- * @param executor thread execution service for consumers on the connection
- * @param addrs an array of known broker addresses (hostname/port pairs) to try in order
- * @return an interface to the connection
- * @throws IOException if it encounters a problem
- */
- public Connection newConnection(ExecutorService executor, Address[] addrs)
- throws IOException
- {
- IOException lastException = null;
- for (Address addr : addrs) {
- try {
- FrameHandler frameHandler = createFrameHandler(addr);
- AMQConnection conn =
- new AMQConnection(username,
- password,
- frameHandler,
- executor,
- virtualHost,
- getClientProperties(),
- requestedFrameMax,
- requestedChannelMax,
- requestedHeartbeat,
- saslConfig);
- conn.start();
- return conn;
- } catch (IOException e) {
- lastException = e;
- }
- }
- throw (lastException != null) ? lastException
- : new IOException("failed to connect");
- }
- /**
- * Create a new broker connection
- * @return an interface to the connection
- * @throws IOException if it encounters a problem
- */
- public Connection newConnection() throws IOException {
- return newConnection(null,
- new Address[] {new Address(getHost(), getPort())}
- );
- }
- /**
- * Create a new broker connection
- * @param executor thread execution service for consumers on the connection
- * @return an interface to the connection
- * @throws IOException if it encounters a problem
- */
- public Connection newConnection(ExecutorService executor) throws IOException {
- return newConnection(executor,
- new Address[] {new Address(getHost(), getPort())}
- );
- }
代码清单2 连接启动
- /**
- * Start up the connection, including the MainLoop thread.
- * Sends the protocol
- * version negotiation header, and runs through
- * Connection.Start/.StartOk, Connection.Tune/.TuneOk, and then
- * calls Connection.Open and waits for the OpenOk. Sets heart-beat
- * and frame max values after tuning has taken place.
- * @throws IOException if an error is encountered
- * either before, or during, protocol negotiation;
- * sub-classes {@link ProtocolVersionMismatchException} and
- * {@link PossibleAuthenticationFailureException} will be thrown in the
- * corresponding circumstances. If an exception is thrown, connection
- * resources allocated can all be garbage collected when the connection
- * object is no longer referenced.
- */
- public void start()
- throws IOException
- {
- this._running = true;
- // Make sure that the first thing we do is to send the header,
- // which should cause any socket errors to show up for us, rather
- // than risking them pop out in the MainLoop
- AMQChannel.SimpleBlockingRpcContinuation connStartBlocker =
- new AMQChannel.SimpleBlockingRpcContinuation();
- // We enqueue an RPC continuation here without sending an RPC
- // request, since the protocol specifies that after sending
- // the version negotiation header, the client (connection
- // initiator) is to wait for a connection.start method to
- // arrive.
- _channel0.enqueueRpc(connStartBlocker);
- try {
- // The following two lines are akin to AMQChannel's
- // transmit() method for this pseudo-RPC.
- _frameHandler.setTimeout(HANDSHAKE_TIMEOUT);
- _frameHandler.sendHeader();
- } catch (IOException ioe) {
- _frameHandler.close();
- throw ioe;
- }
- // start the main loop going
- new MainLoop("AMQP Connection " + getHostAddress() + ":" + getPort()).start();
- // after this point clear-up of MainLoop is triggered by closing the frameHandler.
- AMQP.Connection.Start connStart = null;
- AMQP.Connection.Tune connTune = null;
- try {
- connStart =
- (AMQP.Connection.Start) connStartBlocker.getReply().getMethod();
- _serverProperties = Collections.unmodifiableMap(connStart.getServerProperties());
- Version serverVersion =
- new Version(connStart.getVersionMajor(),
- connStart.getVersionMinor());
- if (!Version.checkVersion(clientVersion, serverVersion)) {
- throw new ProtocolVersionMismatchException(clientVersion,
- serverVersion);
- }
- String[] mechanisms = connStart.getMechanisms().toString().split(" ");
- SaslMechanism sm = this.saslConfig.getSaslMechanism(mechanisms);
- if (sm == null) {
- throw new IOException("No compatible authentication mechanism found - " +
- "server offered [" + connStart.getMechanisms() + "]");
- }
- LongString challenge = null;
- LongString response = sm.handleChallenge(null, this.username, this.password);
- do {
- Method method = (challenge == null)
- ? new AMQP.Connection.StartOk.Builder()
- .clientProperties(_clientProperties)
- .mechanism(sm.getName())
- .response(response)
- .build()
- : new AMQP.Connection.SecureOk.Builder().response(response).build();
- try {
- Method serverResponse = _channel0.rpc(method).getMethod();
- if (serverResponse instanceof AMQP.Connection.Tune) {
- connTune = (AMQP.Connection.Tune) serverResponse;
- } else {
- challenge = ((AMQP.Connection.Secure) serverResponse).getChallenge();
- response = sm.handleChallenge(challenge, this.username, this.password);
- }
- } catch (ShutdownSignalException e) {
- throw new PossibleAuthenticationFailureException(e);
- }
- } while (connTune == null);
- } catch (ShutdownSignalException sse) {
- _frameHandler.close();
- throw AMQChannel.wrap(sse);
- } catch(IOException ioe) {
- _frameHandler.close();
- throw ioe;
- }
- try {
- int channelMax =
- negotiatedMaxValue(this.requestedChannelMax,
- connTune.getChannelMax());
- _channelManager = new ChannelManager(this._workService, channelMax);
- int frameMax =
- negotiatedMaxValue(this.requestedFrameMax,
- connTune.getFrameMax());
- this._frameMax = frameMax;
- int heartbeat =
- negotiatedMaxValue(this.requestedHeartbeat,
- connTune.getHeartbeat());
- setHeartbeat(heartbeat);
- _channel0.transmit(new AMQP.Connection.TuneOk.Builder()
- .channelMax(channelMax)
- .frameMax(frameMax)
- .heartbeat(heartbeat)
- .build());
- _channel0.exnWrappingRpc(new AMQP.Connection.Open.Builder()
- .virtualHost(_virtualHost)
- .build());
- } catch (IOException ioe) {
- _heartbeatSender.shutdown();
- _frameHandler.close();
- throw ioe;
- } catch (ShutdownSignalException sse) {
- _heartbeatSender.shutdown();
- _frameHandler.close();
- throw AMQChannel.wrap(sse);
- }
- // We can now respond to errors having finished tailoring the connection
- this._inConnectionNegotiation = false;
- return;
- }
转载:http://wubin850219.iteye.com/blog/1007984
转载于:https://www.cnblogs.com/telwanggs/p/7124719.html
RabbitMQ学习之ConntectionFactory与Conntection的认知相关推荐
- RabbitMQ学习系列二:.net 环境下 C#代码使用 RabbitMQ 消息队列
上一篇已经讲了Rabbitmq如何在Windows平台安装,不懂请移步:RabbitMQ学习系列一:windows下安装RabbitMQ服务 一.理论: .net环境下,C#代码调用RabbitMQ消 ...
- RabbitMQ学习总结 第一篇:理论篇
目录 RabbitMQ学习总结 第一篇:理论篇 RabbitMQ学习总结 第二篇:快速入门HelloWorld RabbitMQ学习总结 第三篇:工作队列Work Queue RabbitMQ学习总结 ...
- RabbitMQ学习笔记四:RabbitMQ命令(附疑难问题解决)
RabbitMQ学习笔记四:RabbitMQ命令(附疑难问题解决) 参考文章: (1)RabbitMQ学习笔记四:RabbitMQ命令(附疑难问题解决) (2)https://www.cnblogs. ...
- rabbitmq 学习-9- RpcClient发送消息和同步接收消息原理
rabbitmq 学习-9- RpcClient发送消息和同步接收消息原理
- RabbitMQ学习之集群消息可靠性测试
之前介绍过关于消息发送和接收的可靠性:RabbitMQ学习之消息可靠性及特性 下面主要介绍一下集群环境下,rabbitmq实例宕机的情况下,消息的可靠性.验证rabbitmq版本[3.4.1]. ...
- 官网英文版学习——RabbitMQ学习笔记(二)RabbitMQ安装
一.安装RabbitMQ的依赖Erlang 要进行RabbitMQ学习,首先需要进行RabbitMQ服务的安装,安装我们可以根据官网指导进行http://www.rabbitmq.com/downlo ...
- RabbitMQ学习笔记(3)----RabbitMQ Worker的使用
1. Woker队列结构图 这里表示一个生产者生产了消息发送到队列中,但是确有两个消费者在消费同一个队列中的消息. 2. 创建一个生产者 Producer如下: package com.wangx.r ...
- RabbitMQ学习
RabbitMQ学习 1.概述 用于进程通信的中间件. 优势: 劣势: 1.应用解耦:提高了系统容错性和可维护性 1.系统依赖越多不能保证MQ的高可用 2.异步提速:提升用户体验和系统吞吐量 2.复杂 ...
- RabbitMQ 学习笔记
RabbitMQ 学习笔记 RabbitMQ 学习笔记 1. 中间件 1.1 什么是中间件 1.2 为什么要使用消息中间件 1.3 中间件特点 1.4 在项目中什么时候使用中间件技术 2. 中间件技术 ...
最新文章
- 【Android TV 开发】焦点处理 ( 父容器与子组件焦点获取关系处理 | 不同电视设备上的兼容问题 | 触摸获取焦点 | 按键获取焦点 )
- python抢票代码_教你用Python动刷新抢12306火车票,附源码!
- python创建多个txt文件-python-在目录中创建多个文本文件的字数字...
- chrome session丢失_一文带你彻底读懂Cookie、Session、Token到底是什么
- SAP 电商云 Spartacus UI B2B checkout 点击 Continue 不能跳转到下一页面
- 联合体(union)和结构体(struct)的区别
- Multisim破解教程
- 苹果8a1660是什么版本_苹果a1780是什么版本
- 【电路设计】尖峰电压与浪涌电流
- 关于计算机知识的动画电影,动画概论总复习题目(附答案)
- 前端开发公众号的调试
- 基于51单片机智能电子密码锁的设计
- python人脸识别opencv_python中使用Opencv进行人脸识别
- js如何转换json字符串,js如何转换为数值型
- Hexo个人博客NexT主题添加Local Search本地搜索
- 打开i信服务器正在运行中,【网络异常,0/12157 Unknown】i信登录时出现
- 吴军:AI应该变成通识教育,区块链不是炒概念
- LNMP架构部署及应用
- 日语能力考N1复习1
- 帧同步游戏 服务器相关实现
热门文章
- nginx简介--理解nginx配置/模块/openresty
- (32)FPGA面试技能提升篇(EMMC)
- (12)verilog语言编写8路选择器
- linux下dds软件,【数据库】Linux 单实例环境下实现Oracle数据库和DDS软件的开机自动重启...
- QT5日志功能(qDebug、qWarnng、qCritical、qFatal)
- python读取多行json_如何在Python中读取包含多个JSON对象的JSON文件?
- 【蓝桥杯单片机】IIC通讯协议与EEPROM(AT24C02)(官方驱动源码改写)
- 【STM32】【STM32CubeMX】STM32CubeMX的使用之七:定时器输入捕获实现超声波测距
- Linux下pthread的读写锁的优先级问题
- 【MyBatis-Plus】第一章 快速入门