☰
Java Socket多线程聊天室实战:从通信到数据库落地
2026/9/29 19:03:41 网站建设 项目流程

简介:这是一份面向Java初学者与进阶开发者的网络聊天室项目源码,围绕Socket通信、多线程与数据库技术展开,适合用于课程设计、毕业设计或网络编程练手。压缩包共46个文件,约199KB,包含7个java源文件、20个class编译文件、10张jpg与2张gif界面素材、3个txt说明文档,以及db数据库文件、classpath与project工程配置,覆盖源码、资源与工程结构。已有193人学习下载。项目完整呈现了服务端与客户端的Socket连接、消息广播、多线程并发处理、用户注册登录验证及聊天记录存储等核心模块,读者可据此理解网络通信与并发编程的协作方式,并参考其目录组织与代码结构进行二次开发或功能扩展,是掌握Java网络编程综合实践的实用参考。

1. 从一次线上聊天室卡死说起:Java Socket 多线程与数据库到底怎么配合

很多人第一次做 Java 网络聊天室,代码跑在本地两台机器上,消息秒到,觉得自己已经掌握了 socket 网络编程。结果一放到几十人同时在线,服务端线程数飙升、消息开始乱序、数据库连接池直接被打满,控制台刷出一片java.sql.SQLException: IO 错误: socket read timed out。这不是玄学,是典型的「通信层、并发层、持久层」三件事没拆开设计。

这个标题讲的就是一套完整的服务端骨架:用 Java 原生 Socket 做长连接通信,用多线程处理并发客户端,用数据库落地用户、好友、消息记录。它解决的是「消息怎么实时到达、并发怎么不互相阻塞、历史数据怎么不丢」这三个问题。适合已经会写 Hello World 级 Socket、但一上并发就翻车的后端初学者和面试准备者,也适合想用最小依赖手写一遍通信框架、理解 Netty 之前底层长什么样的工程师。下面按「先跑通最小闭环,再补并发和持久化,最后排坑」的顺序讲,每一步都能直接抄。

2. 最小可跑通的 Socket 聊天室:从单连接回显到多客户端广播

先把通信骨架立起来,不碰数据库、不碰线程池,只用最朴素的ServerSocket加一个客户端连接,确认字节流收发没问题。这一步的目的是把「协议格式」定死,后面所有并发和持久化都建立在这个格式上。

2.1 定协议:一行一条消息,字段用竖线分隔

聊天室最容易翻车的地方不是 Socket API,而是消息边界。TCP 是字节流,没有消息概念,你发两次write,对端可能一次read全收到,也可能分三次收到。常见做法是定一个简单文本协议:每条消息以\n结尾,字段用|分隔。

// 协议常量:所有消息统一格式,避免魔法字符串散落各处 public class Protocol { public static final String SEP = "|"; public static final String LOGIN = "LOGIN"; // LOGIN|用户名 public static final String MSG = "MSG"; // MSG|发送者|内容 public static final String SYS = "SYS"; // SYS|系统提示 public static final String QUIT = "QUIT"; // QUIT|用户名 // 组装一条协议消息,末尾必须带换行,作为消息边界 public static String build(String type, String... parts) { StringBuilder sb = new StringBuilder(type); for (String p : parts) { sb.append(SEP).append(p); } return sb.append("\n").toString(); } }

逻辑说明:build方法把类型和参数拼成一行,末尾强制加\n。参数说明:type是消息类型常量,parts是可变参数,顺序必须和解析端约定一致。这样做的代价是内容里不能出现|和换行,真实项目里要么转义,要么换成 JSON 加长度前缀。新手阶段先用这个格式把链路跑通,别一上来就上 Protobuf。

2.2 服务端:一个连接一个线程,先不做线程池

public class ChatServer { // 保存所有在线客户端的输出流,广播时遍历 private static final List<PrintWriter> CLIENTS = Collections.synchronizedList(new ArrayList<>()); public static void main(String[] args) throws IOException { ServerSocket server = new ServerSocket(8888); System.out.println("聊天室启动,端口 8888"); while (true) { Socket socket = server.accept(); // 阻塞等待新连接 new Thread(new ClientHandler(socket)).start(); // 每个连接开一个线程 } } // 处理单个客户端的读写 static class ClientHandler implements Runnable { private final Socket socket; private PrintWriter out; private String name; ClientHandler(Socket socket) { this.socket = socket; } public void run() { try { BufferedReader in = new BufferedReader( new InputStreamReader(socket.getInputStream(), "UTF-8")); out = new PrintWriter(new OutputStreamWriter( socket.getOutputStream(), "UTF-8"), true); CLIENTS.add(out); String line; while ((line = in.readLine()) != null) { // 按行读,靠 \n 切分 String[] p = line.split("\\|"); if (Protocol.LOGIN.equals(p[0])) { name = p[1]; broadcast(Protocol.build(Protocol.SYS, name + " 上线了")); } else if (Protocol.MSG.equals(p[0])) { broadcast(Protocol.build(Protocol.MSG, name, p[2])); } } } catch (IOException e) { System.out.println("连接异常: " + e.getMessage()); } finally { CLIENTS.remove(out); if (name != null) { broadcast(Protocol.build(Protocol.SYS, name + " 下线了")); } try { socket.close(); } catch (IOException ignored) {} } } } // 广播给所有在线客户端 static void broadcast(String msg) { synchronized (CLIENTS) { for (PrintWriter w : CLIENTS) { w.println(msg); // println 自带换行,和协议边界一致 } } } }

逻辑说明:accept()每来一个连接就起一个线程,readLine()按\n切分消息,broadcast遍历所有输出流写回。参数说明:端口 8888 可改,CLIENTS用synchronizedList保证增删线程安全,广播时再手动synchronized一次,防止遍历时被其他线程修改导致ConcurrentModificationException。这一步跑通后,用telnet 127.0.0.1 8888就能手动发消息验证。

2.3 客户端:独立读线程,避免主线程被阻塞

public class ChatClient { public static void main(String[] args) throws IOException { Socket socket = new Socket("127.0.0.1", 8888); PrintWriter out = new PrintWriter(new OutputStreamWriter( socket.getOutputStream(), "UTF-8"), true); // 读线程:专门接收服务端推送,不阻塞下面的键盘输入 new Thread(() -> { try { BufferedReader in = new BufferedReader( new InputStreamReader(socket.getInputStream(), "UTF-8")); String line; while ((line = in.readLine()) != null) { System.out.println(line.replace("|", " ")); } } catch (IOException e) { /* 连接关闭 */ } }).start(); Scanner sc = new Scanner(System.in); System.out.print("输入用户名: "); out.println(Protocol.build(Protocol.LOGIN, sc.nextLine())); while (sc.hasNextLine()) { out.println(Protocol.build(Protocol.MSG, "我", sc.nextLine())); } } }

逻辑说明:客户端必须把「收」和「发」拆到两个线程,否则readLine()会一直阻塞,用户根本没法输入。参数说明:PrintWriter的第二个参数true表示自动 flush,少了它消息会卡在缓冲区发不出去,这是新手最常见的「消息发了但对方收不到」原因。

3. 多线程并发:从裸线程到线程池,把连接数和线程数解耦

第 2 章的「一连接一线程」在几十人时没问题,上百人就开始吃内存,每个线程默认栈 1MB,1000 个连接就是 1GB 栈空间,还没算上下文切换开销。这一章解决并发模型选型,同时把广播的线程安全问题彻底讲清。

3.1 为什么不能无限开线程,线程池参数怎么定

裸new Thread()的问题有三个:创建销毁开销大、数量不可控、异常无法统一处理。常见做法是换成ThreadPoolExecutor,但聊天室是长连接场景,任务不会结束,所以线程池的「任务队列」几乎用不上,核心线程数就等于最大并发连接数。

// 长连接场景:核心线程=最大线程,队列只做缓冲,避免任务被无限堆积 int cores = Runtime.getRuntime().availableProcessors(); ThreadPoolExecutor pool = new ThreadPoolExecutor( cores * 2, // 核心线程数 cores * 2, // 最大线程数,与核心一致 0L, TimeUnit.MILLISECONDS, // 长连接不回收空闲线程 new LinkedBlockingQueue<>(200), // 缓冲队列,防止瞬时连接洪峰 new ThreadFactory() { // 自定义线程名,方便 jstack 排查 private final AtomicInteger n = new AtomicInteger(1); public Thread newThread(Runnable r) { return new Thread(r, "chat-worker-" + n.getAndIncrement()); } }, new ThreadPoolExecutor.AbortPolicy() // 队列满直接拒绝,快速失败 );

逻辑说明:核心线程数和最大线程数设成一样,是因为长连接任务不会释放线程,设大了也没用。参数说明:cores * 2是经验值,IO 密集型可以再高,但聊天室瓶颈通常在网络不在 CPU,2 倍足够;队列 200 是防止瞬间大量连接把内存打爆;AbortPolicy让超载时立刻抛异常,比默默排队更容易发现问题。把第 2 章的new Thread(...).start()换成pool.execute(new ClientHandler(socket))即可。

3.2 广播的线程安全:别在遍历时改集合

第 2 章用了synchronizedList加手动synchronized,能跑但性能差,广播时所有写操作都被阻塞。更稳的做法是换成CopyOnWriteArrayList,读多写少的场景下遍历不加锁。

// 读多写少:广播是高频读,上下线是低频写,CopyOnWrite 最合适 private static final List<PrintWriter> CLIENTS = new CopyOnWriteArrayList<>(); static void broadcast(String msg) { for (PrintWriter w : CLIENTS) { // 遍历的是快照,无需加锁 if (w.checkError()) { // 检测连接是否已断开 CLIENTS.remove(w); // 移除失效连接,避免无效写入 continue; } w.println(msg); } }

逻辑说明:CopyOnWriteArrayList每次写操作复制整个数组,读操作无锁,适合广播这种读远多于写的场景。参数说明:checkError()是PrintWriter自带的错误检测,返回 true 说明底层流已断,及时移除能防止内存泄漏。注意PrintWriter会吞掉 IOException,不检查的话断开的连接会一直留在列表里。

3.3 消息顺序性:多线程下为什么消息会乱序

热搜里常出现「kafka 消费端多线程如何保证消息顺序性」,聊天室同样有这个问题。如果每个客户端消息交给线程池里的不同线程处理,同一个用户连发两条消息,可能第二条先广播出去。解决办法是按会话维度串行:同一个连接的消息始终由同一个线程处理。

// 每个连接绑定一个独立任务,任务内部串行处理该连接的所有消息 // 线程池只负责「连接级」并发,不负责「消息级」并发 public void run() { // 这个 run 方法本身就在一个 worker 线程里跑 // 循环内 readLine 是串行的,天然保证单连接消息有序 while ((line = in.readLine()) != null) { handle(line); // 同步处理,不丢给其他线程 } }

逻辑说明:顺序性的关键是「同一来源的消息不跨线程」。第 2 章的模型天然满足这点,因为一个连接的while循环在一个线程里。参数说明:如果为了吞吐把handle再提交给别的线程池,就必须按用户 ID 做哈希取模路由到固定线程,否则顺序必乱。这是很多人加线程池后消息乱序的血泪经验。

4. 数据库落地:用户、消息、在线状态三张表怎么设计

通信跑通、并发稳住之后,数据不能只存在内存里,重启就没了。这一章把用户信息、消息记录、在线状态落到数据库,重点讲连接池和批量写入,避免「每条消息一个连接」把数据库打挂。

4.1 表结构:三张表覆盖核心场景

表名关键字段用途索引建议
t_userid, username, password_hash, created_at用户账号username 唯一索引
t_messageid, sender, content, send_time消息历史(sender, send_time) 联合索引
t_onlineuser_id, login_time, last_heartbeat在线状态user_id 主键

逻辑说明:t_message只存历史,实时推送走内存广播,两者分离。参数说明:password_hash存哈希不存明文;send_time用datetime而非时间戳,方便直接看;t_online的last_heartbeat用于超时下线判断,客户端每 30 秒发一次心跳更新。

4.2 连接池:为什么不能每次操作都 DriverManager.getConnection

// 用 HikariCP,连接池是数据库访问的生命线 HikariConfig config = new HikariConfig(); config.setJdbcUrl("jdbc:mysql://127.0.0.1:3306/chat?useUnicode=true&characterEncoding=utf8"); config.setUsername("root"); config.setPassword("your_password"); config.setMaximumPoolSize(10); // 最大连接数,别超过数据库 max_connections config.setMinimumIdle(2); // 最小空闲,避免频繁创建 config.setConnectionTimeout(3000); // 获取连接超时 3 秒,快速失败 config.setIdleTimeout(60000); // 空闲连接 60 秒回收 HikariDataSource ds = new HikariDataSource(config);

逻辑说明:连接池复用物理连接,避免每次 TCP 握手和认证。参数说明:maximumPoolSize设 10 是因为聊天室写库不频繁,设太大反而拖垮数据库;connectionTimeout必须设,否则数据库挂了线程会一直等,这就是socket read timed out的常见来源。热搜里的「mysql 的数据库连接池」讲的就是这个,别自己手写池。

4.3 消息落库:批量写入降低数据库压力

// 消息先入内存队列,后台线程批量刷库,避免每条消息一次 insert private final BlockingQueue<Message> queue = new LinkedBlockingQueue<>(10000); // 后台刷库线程 new Thread(() -> { List<Message> batch = new ArrayList<>(100); while (true) { try { Message m = queue.poll(1, TimeUnit.SECONDS); // 最多等 1 秒 if (m != null) batch.add(m); if (batch.size() >= 100 || (m == null && !batch.isEmpty())) { saveBatch(batch); // 批量插入 batch.clear(); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }).start(); // 批量插入,用 addBatch 减少网络往返 void saveBatch(List<Message> list) { String sql = "INSERT INTO t_message(sender, content, send_time) VALUES(?,?,?)"; try (Connection c = ds.getConnection(); PreparedStatement ps = c.prepareStatement(sql)) { for (Message m : list) { ps.setString(1, m.sender); ps.setString(2, m.content); ps.setTimestamp(3, new Timestamp(m.time)); ps.addBatch(); } ps.executeBatch(); } catch (SQLException e) { System.out.println("批量落库失败: " + e.getMessage()); } }

逻辑说明:消息先进队列,后台线程攒够 100 条或超时 1 秒就批量写,把 N 次网络往返压成 1 次。参数说明:队列容量 10000 是防止内存无限增长,满了就丢最老的消息或阻塞;批量大小 100 是吞吐和延迟的折中,太大延迟高,太小没效果。注意executeBatch失败要记录,别静默吞掉。

5. 避坑与排查:那些让聊天室半夜挂掉的细节

这一章全是踩过的坑,每条按「现象 → 原因 → 解决」写,遇到问题直接对号入座。

5.1 现象:客户端发了消息,服务端收不到

原因:PrintWriter没开自动 flush,或者用了BufferedWriter忘了flush()。TCP 有缓冲区,数据攒着不发是常态。解决:PrintWriter构造时第二个参数传true,或者每次写完手动flush()。这是新手第一大坑,占「消息发不出去」问题的一半以上。

5.2 现象:中文乱码,英文正常

原因:InputStreamReader和OutputStreamWriter没指定字符集,用了平台默认编码,Windows 是 GBK,Linux 是 UTF-8,两端不一致就乱码。解决:两端都显式指定"UTF-8",数据库连接串也加characterEncoding=utf8。别依赖默认值,跨平台必翻车。

5.3 现象:连接数一多就报Too many connections

原因:每个连接一个数据库连接,或者连接用完没关闭。解决:用连接池,try-with-resources保证Connection、Statement、ResultSet自动关闭。检查代码里有没有getConnection()后忘了close()的分支,尤其是异常路径。

5.4 现象:服务端线程数只增不减,内存持续上涨

原因:客户端异常断开时,readLine()抛异常,但线程没退出,或者CLIENTS列表里的失效PrintWriter没移除。解决:finally块里必须CLIENTS.remove(out)和socket.close();广播时用checkError()清理失效连接。用jstack看线程名,如果chat-worker-数量只增不减,就是这个原因。

5.5 现象:消息偶尔丢失或重复

原因:广播时遍历CLIENTS被其他线程修改,或者批量落库时队列满了丢消息。解决:集合用CopyOnWriteArrayList;队列满时要么阻塞生产者,要么记录丢弃日志,别静默丢。消息重复通常是客户端重连后重发,需要消息 ID 去重,这个属于进阶话题。

6. 进阶技巧:用心跳检测和优雅关闭把服务端做扎实

前面五章跑通了功能,但一个能长期运行的服务端还得处理「死连接」和「优雅停机」。这一章讲两个具体技巧,都是线上环境必须的。

6.1 心跳检测:怎么判断客户端是真在线还是假死

TCP 连接断开时,如果客户端是拔网线或进程崩溃,服务端可能长时间不知道,readLine()一直阻塞。解决办法是心跳:客户端每 30 秒发一条PING,服务端更新t_online.last_heartbeat,后台线程扫描超过 90 秒没心跳的连接主动关闭。

// 服务端心跳扫描线程,每 30 秒检查一次 ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(() -> { long now = System.currentTimeMillis(); for (ClientHandler h : HANDLERS) { // 维护一个 handler 集合 if (now - h.lastActive > 90_000) { // 90 秒无活动 h.close(); // 主动关闭,触发 finally 清理 } } }, 30, 30, TimeUnit.SECONDS);

逻辑说明:lastActive在每次收到消息时更新,扫描线程只读不写,关闭操作交给 handler 自己的close()。参数说明:心跳间隔 30 秒、超时 90 秒是经验值,间隔太短浪费流量,太长发现死连接慢。注意close()要幂等,重复调用不能抛异常。

6.2 优雅关闭:停机时别丢消息

直接kill -9会丢内存队列里没落库的消息。正确做法是注册 JVM 关闭钩子,先停止接收新连接,再把队列刷完,最后关连接池。

Runtime.getRuntime().addShutdownHook(new Thread(() -> { System.out.println("开始优雅关闭..."); server.close(); // 1. 停止 accept 新连接 pool.shutdown(); // 2. 等待线程池任务结束 try { pool.awaitTermination(5, TimeUnit.SECONDS); } catch (InterruptedException ignored) {} flushQueue(); // 3. 把内存队列剩余消息落库 ds.close(); // 4. 关闭连接池 System.out.println("关闭完成"); }));

逻辑说明:关闭钩子在 JVM 收到SIGTERM时执行,顺序不能乱,先停新连接再刷数据。参数说明:awaitTermination等 5 秒是给正在处理的消息留时间,超时就强制继续。这套流程在容器环境里尤其重要,K8s 滚动更新时会发SIGTERM,没有钩子就会丢消息。

我自己做这类项目最大的教训是:别急着上框架。先把原生 Socket 加线程池加连接池这套最小闭环手写一遍,把消息边界、线程安全、连接生命周期这三个点吃透,后面换 Netty 或 Spring Boot 的 WebSocket 时,才知道每个配置项背后在解决什么问题。很多人直接抄框架代码,一出问题就抓瞎,就是因为底层黑匣子没打开过。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询