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

如何为线程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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:51:13