数据库路由中间件MyCat - 源代码篇(4)
2. 前端连接建立与认证
Created with Raphaël 2.1.0 MySql连接建立以及认证过程 client client MySql MySql 1.TCP连接请求 2.接受TCP连接 3.TCP连接建立 4.握手包HandshakePacket 5.认证包AuthPacket 6.如果验证成功,则返回OkPacket 7.默认会发送查询版本信息的包 8.返回结果包
2.5 (7~8) 默认会发送查询版本信息的包,返回结果包
MySql客户端在连接建立后,默认会发送查询版本信息的包,这其实就是一个SQL查询请求了。只不过这个请求不用路由到后台某个数据库^_^。
连接成功建立后,连接绑定的RW线程会监听上面的读事件。在客户端发送查询版本信息的包之后,会触发RW线程去读取对应连接,过程与之前接收AuthPacket类似:
RW类代码片段
1//监听到有效读 2if (key.isValid() && key.isReadable()) { 3 try { 4 //异步读取数据并处理数据 5 con.asynRead(); 6 } catch (IOException e) { 7 con.close("program err:" + e.toString()); 8 continue; 9 } catch (Exception e) { 10 LOGGER.debug("caught err:", e); 11 con.close("program err:" + e.toString()); 12 continue; 13 } 14 }
之后的读取过程也是调用AbstractConnection的asynRead()方法,进行异步读取。过程就不再赘述,读取到的数据交由FrontendCommandHandler处理。
查询版本信息的包(是一种CommandPacket)内容:

CommandPacket:
- packet length (3)
- packet number (1)
- command (1)
- statement (null terminated string)
FrontendCommandHandler的处理方法:
1@Override 2 public void handle(byte[] data) 3 { 4 5 if(source.getLoadDataInfileHandler()!=null&&source.getLoadDataInfileHandler().isStartLoadData()) 6 { 7 MySQLMessage mm = new MySQLMessage(data); 8 int packetLength = mm.readUB3(); 9 if(packetLength+4==data.length) 10 { 11 source.loadDataInfileData(data); 12 } 13 return; 14 } 15 switch (data[4]) 16 { 17 case MySQLPacket.COM_INIT_DB: 18 commands.doInitDB(); 19 source.initDB(data); 20 break; 21 case MySQLPacket.COM_QUERY: 22 commands.doQuery(); 23 source.query(data); 24 break; 25 case MySQLPacket.COM_PING: 26 commands.doPing(); 27 source.ping(); 28 break; 29 case MySQLPacket.COM_QUIT: 30 commands.doQuit(); 31 source.close("quit cmd"); 32 break; 33 case MySQLPacket.COM_PROCESS_KILL: 34 commands.doKill(); 35 source.kill(data); 36 break; 37 case MySQLPacket.COM_STMT_PREPARE: 38 commands.doStmtPrepare(); 39 source.stmtPrepare(data); 40 break; 41 case MySQLPacket.COM_STMT_EXECUTE: 42 commands.doStmtExecute(); 43 source.stmtExecute(data); 44 break; 45 case MySQLPacket.COM_STMT_CLOSE: 46 commands.doStmtClose(); 47 source.stmtClose(data); 48 break; 49 case MySQLPacket.COM_HEARTBEAT: 50 commands.doHeartbeat(); 51 source.heartbeat(data); 52 break; 53 default: 54 commands.doOther(); 55 source.writeErrMessage(ErrorCode.ER_UNKNOWN_COM_ERROR, 56 "Unknown command"); 57 58 } 59 }
根据CommandPacket的第五字节判断command类型,不同类型有不同的处理。
首先querycommand计数加1,之后调用对应FrontendConnection的query(byte[])方法:
1public void query(byte[] data) { 2 if (queryHandler != null) { 3 // 取得语句|get sql 4 MySQLMessage mm = new MySQLMessage(data); 5 //从第六字节开始读取|read from the 6th byte 6 mm.position(5); 7 String sql = null; 8 try { 9 sql = mm.readString(charset); 10 } catch (UnsupportedEncodingException e) { 11 writeErrMessage(ErrorCode.ER_UNKNOWN_CHARACTER_SET, "Unknown charset '" + charset + "'"); 12 return; 13 } 14 if (sql == null || sql.length() == 0) { 15 writeErrMessage(ErrorCode.ER_NOT_ALLOWED_COMMAND, "Empty SQL"); 16 return; 17 } 18 19 // sql = StringUtil.replace(sql, "`", ""); 20 21 // 移除末尾';'|remove last ';' 22 if (sql.endsWith(";")) { 23 sql = sql.substring(0, sql.length() - 1); 24 } 25 26 // 记录SQL|record SQL 27 this.setExecuteSql(sql); 28 29 // 执行查询 30 queryHandler.setReadOnly(privileges.isReadOnly(user)); 31 queryHandler.query(sql); 32 } else { 33 writeErrMessage(ErrorCode.ER_UNKNOWN_COM_ERROR, "Query unsupported!"); 34 } 35 }
执行查询,调用对应的FrontendQueryHandler:

这里,很明显,是ServerQueryHandler。
1public void query(String sql) { 2 3 ServerConnection c = this.source; 4 if (LOGGER.isDebugEnabled()) { 5 LOGGER.debug(new StringBuilder().append(c).append(sql).toString()); 6 } 7 // 8 int rs = ServerParse.parse(sql); 9 int sqlType = rs & 0xff; 10 11 switch (sqlType) { 12 case ServerParse.EXPLAIN: 13 ExplainHandler.handle(sql, c, rs >>> 8); 14 break; 15 case ServerParse.EXPLAIN2: 16 Explain2Handler.handle(sql, c, rs >>> 8); 17 break; 18 case ServerParse.SET: 19 SetHandler.handle(sql, c, rs >>> 8); 20 break; 21 case ServerParse.SHOW: 22 ShowHandler.handle(sql, c, rs >>> 8); 23 break; 24 case ServerParse.SELECT: 25 if(QuarantineHandler.handle(sql, c)){ 26 SelectHandler.handle(sql, c, rs >>> 8); 27 } 28 break; 29 case ServerParse.START: 30 StartHandler.handle(sql, c, rs >>> 8); 31 break; 32 case ServerParse.BEGIN: 33 BeginHandler.handle(sql, c); 34 break; 35 case ServerParse.SAVEPOINT: 36 SavepointHandler.handle(sql, c); 37 break; 38 case ServerParse.KILL: 39 KillHandler.handle(sql, rs >>> 8, c); 40 break; 41 case ServerParse.KILL_QUERY: 42 LOGGER.warn(new StringBuilder().append("Unsupported command:").append(sql).toString()); 43 c.writeErrMessage(ErrorCode.ER_UNKNOWN_COM_ERROR,"Unsupported command"); 44 break; 45 case ServerParse.USE: 46 UseHandler.handle(sql, c, rs >>> 8); 47 break; 48 case ServerParse.COMMIT: 49 c.commit(); 50 break; 51 case ServerParse.ROLLBACK: 52 c.rollback(); 53 break; 54 case ServerParse.HELP: 55 LOGGER.warn(new StringBuilder().append("Unsupported command:").append(sql).toString()); 56 c.writeErrMessage(ErrorCode.ER_SYNTAX_ERROR, "Unsupported command"); 57 break; 58 case ServerParse.MYSQL_CMD_COMMENT: 59 c.write(c.writeToBuffer(OkPacket.OK, c.allocate())); 60 break; 61 case ServerParse.MYSQL_COMMENT: 62 c.write(c.writeToBuffer(OkPacket.OK, c.allocate())); 63 break; 64 case ServerParse.LOAD_DATA_INFILE_SQL: 65 c.loadDataInfileStart(sql); 66 break; 67 default: 68 if(readOnly){ 69 LOGGER.warn(new StringBuilder().append("User readonly:").append(sql).toString()); 70 c.writeErrMessage(ErrorCode.ER_USER_READ_ONLY, "User readonly"); 71 break; 72 } 73 if(QuarantineHandler.handle(sql, c)){ 74 c.execute(sql, rs & 0xff); 75 } 76 } 77 }
针对每种command,都有不同的handler和处理方式。之后如何处理,就在之后的SQL解析器等章节进行分析。