资讯动态

Java Socket多线程聊天室实战:从通信到数据库落地

发布时间:2026/9/29 19:03:48 来源:尧图企业网站定制
简介这是一份面向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 ListPrintWriter 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 章的「一连接一线程」在几十人时没问题上百人就开始吃内存每个线程默认栈 1MB1000 个连接就是 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 密集型可以再高但聊天室瓶颈通常在网络不在 CPU2 倍足够队列 200 是防止瞬间大量连接把内存打爆AbortPolicy让超载时立刻抛异常比默默排队更容易发现问题。把第 2 章的new Thread(...).start()换成pool.execute(new ClientHandler(socket))即可。3.2 广播的线程安全别在遍历时改集合第 2 章用了synchronizedList加手动synchronized能跑但性能差广播时所有写操作都被阻塞。更稳的做法是换成CopyOnWriteArrayList读多写少的场景下遍历不加锁。// 读多写少广播是高频读上下线是低频写CopyOnWrite 最合适 private static final ListPrintWriter 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?useUnicodetruecharacterEncodingutf8); 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 BlockingQueueMessage queue new LinkedBlockingQueue(10000); // 后台刷库线程 new Thread(() - { ListMessage 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(ListMessage 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 是 GBKLinux 是 UTF-8两端不一致就乱码。解决两端都显式指定UTF-8数据库连接串也加characterEncodingutf8。别依赖默认值跨平台必翻车。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 时才知道每个配置项背后在解决什么问题。很多人直接抄框架代码一出问题就抓瞎就是因为底层黑匣子没打开过。希望帮到你。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑