From 0c7aed197a5a4f33cdab13121791e2772b99db70 Mon Sep 17 00:00:00 2001 From: jason huang Date: Tue, 2 Jan 2018 23:29:15 +0800 Subject: [PATCH] cannot build mysql connection. when invoke connect(), it will return a ChannelFuture instance, and then the "addListener" method is invoked by biz thread, however, "channelRead" method in BusinessHandler is executed by netty I/O thread. ChannelFutureListener -> operationComplete() cannot guarantee that which is executed before than BusinessHandler -> channelRead() method. reproduce it : set below jvm parameter: -Dio.netty.eventLoopThreads=1 It will caused new connection is timeout when AbstractEventParser -> start() -> parseThread ->erosaConnection.reconnect() is invoked. --- .../mysql/socket/SocketChannelPool.java | 37 +++++++++++-------- 1 file changed, 21 insertions(+), 16 deletions(-) diff --git a/driver/src/main/java/com/alibaba/otter/canal/parse/driver/mysql/socket/SocketChannelPool.java b/driver/src/main/java/com/alibaba/otter/canal/parse/driver/mysql/socket/SocketChannelPool.java index 52609b19..53b78f06 100644 --- a/driver/src/main/java/com/alibaba/otter/canal/parse/driver/mysql/socket/SocketChannelPool.java +++ b/driver/src/main/java/com/alibaba/otter/canal/parse/driver/mysql/socket/SocketChannelPool.java @@ -54,31 +54,25 @@ public abstract class SocketChannelPool { } public static SocketChannel open(SocketAddress address) throws Exception { - final SocketChannel socket = new SocketChannel(); - final BooleanMutex mutex = new BooleanMutex(false); - boot.connect(address).addListener(new ChannelFutureListener() { + SocketChannel socket = null; + ChannelFuture future = boot.connect(address).sync(); - @Override - public void operationComplete(ChannelFuture arg0) throws Exception { - if (arg0.isSuccess()) { - socket.setChannel(arg0.channel()); - } + if (future.isSuccess()) { + future.channel().pipeline().get(BusinessHandler.class).latch.await(); + socket = chManager.get(future.channel()); + } - mutex.set(true); - } - }); - // wait for complete - mutex.get(); - if (null == socket.getChannel()) { + if (null == socket) { throw new IOException("can't create socket!"); } - chManager.put(socket.getChannel(), socket); + return socket; } public static class BusinessHandler extends ChannelInboundHandlerAdapter { private SocketChannel socket = null; + private final CountDownLatch latch = new CountDownLatch(1); @Override public void channelInactive(ChannelHandlerContext ctx) throws Exception { @@ -86,11 +80,22 @@ public abstract class SocketChannelPool { chManager.remove(ctx.channel());// 移除 } + @Override + public void channelActive(ChannelHandlerContext ctx) throws Exception { + socket = new SocketChannel(); + socket.setChannel(ctx.channel()); + chManager.put(ctx.channel(), socket); + latch.countDown(); + super.channelActive(ctx); + } + @Override public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception { - if (null == socket) socket = chManager.get(ctx.channel()); if (socket != null) { socket.writeCache((ByteBuf) msg); + } else { + //TODO: need graceful error handler. + logger.error("no socket available."); } ReferenceCountUtil.release(msg);// 添加防止内存泄漏的 }