Android中AsynchronousSocketChannel跨线程运行解决主线程网络异常
解决Android中AsynchronousSocketChannel的NetworkOnMainThreadException问题
首先咱们得揪出问题根源:Android主线程(UI线程)严格禁止任何网络操作,哪怕你用的是AsynchronousSocketChannel这种异步API——你在主线程里发起了connect调用,再加上MainActivity里写的无限while循环,不仅会触发NetworkOnMainThreadException,还会直接导致ANR(应用无响应)。
和Boost.Asio用独立线程运行io_context的思路完全一致,我们可以给AsynchronousSocketChannel指定自定义线程池(Executor),让所有网络异步操作的回调都在独立线程里执行,彻底和主线程划清界限。下面是具体解决方案:
1. 给AsynchronousSocketChannel配置自定义Executor
AsynchronousSocketChannel.open()支持传入Executor参数,我们用它指定自己的线程池,这样所有异步操作(connect、read、write)的回调都会在这个线程池的线程里跑,不会占用主线程资源。
2. 重构Ethernet类
修改你的Ethernet类,加入自定义线程池,并修正代码里的小问题:
import java.io.IOException; import java.net.InetSocketAddress; import java.nio.ByteBuffer; import java.nio.channels.AsynchronousSocketChannel; import java.nio.channels.CompletionHandler; import java.util.Queue; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import android.util.Log; public class Ethernet { private String ip; private int port; private Queue<String> recvDataQueue, sendDataQueue; private InetSocketAddress address; private AsynchronousSocketChannel socket; private boolean socketAlive, connectionInProgress; private ByteBuffer readBuffer, sendBuffer; // 自定义网络线程池,对应Boost.Asio的io_context线程 private ExecutorService networkExecutor; Ethernet(String ip, int port) { if (port > 65535 || port < 1024) { // 修正原条件判断逻辑 port = 6666; } this.ip = ip; this.port = port; this.socketAlive = false; connectionInProgress = false; address = new InetSocketAddress(ip, port); this.readBuffer = ByteBuffer.allocate(8192); // 初始化单线程线程池,也可根据需求用FixedThreadPool networkExecutor = Executors.newSingleThreadExecutor(); } public boolean connect() { try { if (this.socketAlive || connectionInProgress) { return false; } connectionInProgress = true; // 打开Channel时传入自定义Executor this.socket = AsynchronousSocketChannel.open(networkExecutor); this.socket.connect(this.address, null, new CompletionHandler<Void, Object>() { @Override public void completed(Void result, Object attachment) { socketAlive = true; connectionInProgress = false; Log.d("Ethernet", "连接成功"); // 连接成功后启动数据监听 try { recieveDataSocket(); } catch (Exception e) { Log.e("Ethernet", "启动接收失败: " + e.getMessage()); } } @Override public void failed(Throwable e, Object attachment) { socketAlive = false; connectionInProgress = false; Log.e("Ethernet", "连接失败: " + e.getMessage()); try { if (socket != null && socket.isOpen()) { socket.close(); } } catch (IOException ex) { Log.e("Ethernet", "关闭socket失败: " + ex.getMessage()); } } }); return true; } catch (Exception e) { Log.e("Ethernet", "连接异常: " + e.getMessage()); connectionInProgress = false; } return false; } public long sendData(String writeData) { if (!socketAlive || socket == null || !socket.isOpen()) { Log.e("Ethernet", "socket未连接,无法发送数据"); return -1; } try { sendBuffer = ByteBuffer.wrap(writeData.getBytes()); socket.write(sendBuffer, null, new CompletionHandler<Integer, Object>() { @Override public void completed(Integer result, Object attachment) { if (result < 0) { // 连接意外关闭 socketAlive = false; connectionInProgress = false; Log.d("Ethernet", "连接意外关闭"); } else if (sendBuffer.hasRemaining()) { // 剩余数据未发送,继续发送 socket.write(sendBuffer, null, this); } else { Log.d("Ethernet", "数据发送完成,共发送" + result + "字节"); sendBuffer.clear(); } } @Override public void failed(Throwable e, Object attachment) { Log.e("Ethernet", "发送数据失败: " + e.getMessage()); if (socket != null && socket.isOpen()) { try { socketAlive = false; socket.close(); } catch (IOException ex) { Log.e("Ethernet", "关闭socket失败: " + ex.getMessage()); } } } }); return writeData.length(); } catch (Exception e) { Log.e("Ethernet", "发送数据异常: " + e.getMessage()); } return -1; } private void recieveDataSocket() throws Exception { if (!socketAlive || socket == null || !socket.isOpen()) { return; } readBuffer.clear(); socket.read(readBuffer, null, new CompletionHandler<Integer, Object>() { @Override public void completed(Integer result, Object attachment) { if (result < 0) { // 连接关闭 socketAlive = false; connectionInProgress = false; Log.d("Ethernet", "连接关闭"); try { if (socket != null && socket.isOpen()) { socket.close(); } } catch (IOException ex) { Log.e("Ethernet", "关闭socket失败: " + ex.getMessage()); } } else { readBuffer.flip(); byte[] data = new byte[readBuffer.remaining()]; readBuffer.get(data); String recvStr = new String(data); Log.d("Ethernet", "收到数据: " + recvStr); // 可将数据加入recvDataQueue // recvDataQueue.add(recvStr); // 继续监听下一次数据 try { recieveDataSocket(); } catch (Exception e) { Log.e("Ethernet", "重启接收失败: " + e.getMessage()); } } } @Override public void failed(Throwable e, Object attachment) { Log.e("Ethernet", "接收数据失败: " + e.getMessage()); if (socketAlive && socket != null && socket.isOpen()) { try { socketAlive = false; socket.close(); } catch (IOException ex) { Log.e("Ethernet", "关闭socket失败: " + ex.getMessage()); } } } }); } public void disconnect() { try { if (socket != null && socket.isOpen()) { socket.close(); } socketAlive = false; connectionInProgress = false; // 关闭线程池释放资源 if (!networkExecutor.isShutdown()) { networkExecutor.shutdown(); } } catch (IOException ex) { Log.e("Ethernet", "断开连接异常: " + ex.getMessage()); } } // 补充原代码缺失的getter方法 public boolean isSocketAlive() { return socketAlive; } public boolean isConnectionInProgress() { return connectionInProgress; } }
3. 修改MainActivity的逻辑
绝对不能在主线程写无限while循环,这会直接卡死UI。我们把连接和发送逻辑放到自定义线程池里执行:
import android.os.Bundle; import androidx.appcompat.app.AppCompatActivity; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; public class MainActivity extends AppCompatActivity { private Ethernet eth; // 执行连接逻辑的线程池,也可复用Ethernet里的networkExecutor private ExecutorService mainNetworkExecutor; @Override protected void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(R.layout.activity_main); eth = new Ethernet("192.168.1.22", 6666); mainNetworkExecutor = Executors.newSingleThreadExecutor(); // 在子线程执行连接和发送逻辑 mainNetworkExecutor.submit(() -> { while (!Thread.currentThread().isInterrupted()) { if (!eth.isSocketAlive() && !eth.isConnectionInProgress()) { Log.d("MainActivity", "尝试连接服务器"); eth.connect(); } if (eth.isSocketAlive()) { Log.d("MainActivity", "发送测试数据"); eth.sendData("Hello, World!!!\n"); // 不要立刻断开,否则会反复连接断开,可根据需求调整 // eth.disconnect(); // 加延迟避免频繁发送 try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } try { Thread.sleep(500); // 避免循环太频繁占用CPU } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }); } @Override protected void onDestroy() { super.onDestroy(); // 销毁时断开连接、关闭线程池 eth.disconnect(); if (!mainNetworkExecutor.isShutdown()) { mainNetworkExecutor.shutdownNow(); } } }
关键说明
- 自定义Executor就相当于Boost.Asio里
io_context.run()所在的线程,所有异步操作的回调都会在这个线程池的线程中执行,完全避开主线程的网络限制。 - 主线程绝对不能做阻塞操作(比如原代码里的无限while循环),必须把这类逻辑放到子线程或线程池里。
- 记得在Activity销毁时关闭socket和线程池,避免内存泄漏。
内容的提问来源于stack exchange,提问作者Dark Sorrow
相关产品推荐
相关产品推荐

