如何为线程A写Buffer、线程B读存数据库的代码实现线程安全?
线程安全实现:线程A与线程B的交替协作方案
看起来你需要的是线程间的有序交替执行:线程A写完数据后,线程B才能读取并保存;线程B完成保存后,线程A才能继续写下一批数据。这是典型的生产者-消费者模型的变种,核心是线程间的协作与同步,我会用两种方式帮你实现,一种是基于synchronized+wait/notify的手动实现(帮你理解底层原理),另一种是用Java自带的线程安全队列快速实现(更简洁高效)。
一、手动实现:synchronized + wait/notify
我们先构建一个共享的Buffer类,用它来承载数据,同时通过状态标记和线程等待唤醒机制来保证交替执行:
1. 线程安全的共享Buffer类
class SharedBuffer { private String data; // 标记Buffer中是否有可读取的数据 private boolean hasData = false; // 线程A调用:写入数据 public synchronized void write(String newData) throws InterruptedException { // 如果Buffer里已有未读取的数据,线程A进入等待 while (hasData) { wait(); // 释放当前对象锁,进入等待队列 } this.data = newData; hasData = true; System.out.println("线程A写入数据: " + data); notifyAll(); // 唤醒所有等待该锁的线程(这里就是线程B) } // 线程B调用:读取数据 public synchronized String read() throws InterruptedException { // 如果Buffer为空,线程B进入等待 while (!hasData) { wait(); } String result = this.data; hasData = false; System.out.println("线程B读取数据: " + result); notifyAll(); // 唤醒所有等待该锁的线程(这里就是线程A) return result; } }
2. 实现线程A(生产者)和线程B(消费者)
// 线程A:负责写入数据到Buffer class WriterThread extends Thread { private SharedBuffer buffer; private String[] dataToWrite; public WriterThread(SharedBuffer buffer, String[] dataToWrite) { this.buffer = buffer; this.dataToWrite = dataToWrite; } @Override public void run() { try { for (String data : dataToWrite) { buffer.write(data); // 模拟写入数据的耗时操作(比如从文件读取数据) Thread.sleep(500); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("写入线程被中断"); } } } // 线程B:负责读取Buffer数据并保存到数据库 class ReaderThread extends Thread { private SharedBuffer buffer; // 模拟保存到数据库的操作 private void saveToDatabase(String data) { System.out.println("线程B将数据保存到数据库: " + data); } public ReaderThread(SharedBuffer buffer) { this.buffer = buffer; } @Override public void run() { try { // 这里假设要读取的次数和写入次数一致 for (int i = 0; i < 3; i++) { String data = buffer.read(); saveToDatabase(data); // 模拟数据库保存的耗时操作 Thread.sleep(800); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("读取线程被中断"); } } }
3. 主函数启动线程
public class Main { public static void main(String args[]) throws InterruptedException { SharedBuffer buffer = new SharedBuffer(); // 准备要写入的测试数据 String[] testData = {"用户数据1", "用户数据2", "用户数据3"}; Thread threadA = new WriterThread(buffer, testData); Thread threadB = new ReaderThread(buffer); // 启动两个线程 threadA.start(); threadB.start(); // 等待两个线程全部执行完毕 threadA.join(); threadB.join(); System.out.println("所有数据写入与保存操作完成"); } }
关键细节说明
- synchronized关键字:修饰
write和read方法,保证同一时间只有一个线程能访问这些方法,避免并发修改Buffer的状态导致数据不一致。 - wait()与notifyAll():
wait()会让当前线程释放锁并进入等待状态,直到其他线程调用notifyAll()唤醒它。这里用while循环判断状态而不是if,是为了防止虚假唤醒(线程可能在未被正确通知的情况下醒来,需要再次检查状态是否满足执行条件)。 - 中断处理:在捕获
InterruptedException后,调用Thread.currentThread().interrupt()保留中断状态,方便上层代码处理中断逻辑。
二、快速实现:使用LinkedBlockingQueue
如果你不想手动处理线程同步,Java的java.util.concurrent.LinkedBlockingQueue是现成的线程安全阻塞队列,它自带阻塞的put()和take()方法,能天然实现你的需求:
import java.util.concurrent.LinkedBlockingQueue; public class Main { public static void main(String args[]) throws InterruptedException { // 设置队列容量为1,保证每次只能存一个数据,强制线程交替执行 LinkedBlockingQueue<String> queue = new LinkedBlockingQueue<>(1); String[] testData = {"用户数据1", "用户数据2", "用户数据3"}; // 线程A:写入数据到队列 Thread threadA = new Thread(() -> { try { for (String data : testData) { queue.put(data); // 队列满时自动阻塞 System.out.println("线程A写入数据: " + data); Thread.sleep(500); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("写入线程被中断"); } }); // 线程B:读取队列数据并保存到数据库 Thread threadB = new Thread(() -> { try { for (int i = 0; i < testData.length; i++) { String data = queue.take(); // 队列空时自动阻塞 System.out.println("线程B读取数据: " + data); System.out.println("线程B将数据保存到数据库: " + data); Thread.sleep(800); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); System.out.println("读取线程被中断"); } }); threadA.start(); threadB.start(); threadA.join(); threadB.join(); System.out.println("所有操作完成"); } }
方案优势
LinkedBlockingQueue是Java并发包提供的线程安全队列,内部已经实现了synchronized或Lock的同步逻辑,无需手动处理锁和等待唤醒。put()方法在队列满时自动阻塞线程,take()方法在队列空时自动阻塞线程,完美匹配你的“线程A等线程B读完再写,线程B等线程A写完再读”的需求。
内容的提问来源于stack exchange,提问作者Leonardo
相关产品推荐
相关产品推荐

