You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.11 09:28:36