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

Java中如何将消费者循环的数据传递至EmailServices类

解决方案:将EmailBuffer的验证数据传递至EmailServices并优化实现

一、核心修改步骤

1. 修正EmailBuffer类,让get()方法返回EmailModel

原EmailBuffer的泛型声明错误(自定义了<EmailModel>类型参数,而非引用你定义的com.webservices.model.EmailModel),且get()方法未返回取出的数据。修改后让它返回验证后的邮件数据:

package com.webservices.buffer;

import java.util.concurrent.LinkedBlockingQueue;
import com.webservices.model.EmailModel;

//Shared class used by threads
public class EmailBuffer {
    // 移除错误泛型声明,直接使用EmailModel;调整队列容量避免频繁阻塞
    private LinkedBlockingQueue<EmailModel> emailLinkedBlockingQueue = new LinkedBlockingQueue<>(10);

    // 修改get方法,返回取出的EmailModel
    public EmailModel get() throws InterruptedException {
        EmailModel emailData = emailLinkedBlockingQueue.take();
        System.out.println("Consumer received - " + emailData);
        return emailData;
    }

    public void put(EmailModel user) throws InterruptedException {
        emailLinkedBlockingQueue.put(user);
        System.out.println("Producer produced - " + user);
    }
}

2. 修改EmailConsumer类,传递数据至EmailServices

让消费者从缓冲区拿到数据后,提交邮件发送任务到线程池(避免频繁创建线程):

package com.webservices.worker;

import com.webservices.buffer.EmailBuffer;
import com.webservices.services.EmailServices;
import com.webservices.model.EmailModel;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class EmailConsumer implements Runnable {
    private EmailBuffer buffer;
    private ExecutorService emailExecutor; // 线程池管理邮件发送任务

    public EmailConsumer(EmailBuffer buffer) {
        this.buffer = buffer;
        // 根据业务量设置线程池大小
        this.emailExecutor = Executors.newFixedThreadPool(5);
    }
    
    @Override
    public void run() {
        while (!Thread.currentThread().isInterrupted()) { // 用中断实现优雅停止
            try {
                EmailModel emailData = buffer.get();
                // 提交邮件发送任务
                emailExecutor.submit(new EmailServices(emailData));
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt(); // 重置中断状态
                System.out.println("Consumer thread interrupted");
                break;
            }
        }
        emailExecutor.shutdown(); // 关闭线程池
    }
}

3. 调整EmailServices类,接收并使用EmailModel数据

让邮件服务通过构造函数接收验证后的EmailModel,复用JavaMail的Session减少资源消耗:

package com.webservices.services;

import java.util.Properties;
import javax.mail.Message;
import javax.mail.PasswordAuthentication;
import javax.mail.Session;
import javax.mail.Transport;
import javax.mail.internet.InternetAddress;
import javax.mail.internet.MimeMessage;
import com.webservices.model.EmailModel;

public class EmailServices implements Runnable {
    private static Session mailSession; // 复用Session,仅初始化一次
    private EmailModel emailModel;

    // 静态初始化Session
    static {
        String sender = "xyz@gmail.com";
        String smtpHost = "127.0.0.1"; // 修正原无效的127.0.0.0地址
        int smtpPort = 587;
        
        Properties properties = System.getProperties(); 
        properties.put("mail.smtp.host", smtpHost); 
        properties.put("mail.smtp.port", smtpPort); 
        properties.put("mail.smtp.ssl.protocols", "TLSv1.2");
        properties.put("mail.smtp.starttls.enable", "true");
        properties.put("mail.smtp.auth", "true");
        properties.put("mail.debug", "false"); // 生产环境关闭debug
        
        mailSession = Session.getInstance(properties,
             new javax.mail.Authenticator() {
            @Override
             protected PasswordAuthentication getPasswordAuthentication()
             {
                 return new PasswordAuthentication(sender,"*****");
             }
         });
        mailSession.setDebug(false);
    }

    // 构造函数接收EmailModel
    public EmailServices(EmailModel emailModel) {
        this.emailModel = emailModel;
    }

    @Override
    public void run() {
        System.out.println("Preparing to send email to: " + emailModel.getEmail());
        try {
            sendEmail(emailModel.getEmail(), emailModel.getSubject(), emailModel.getMessage());
        } catch (Exception e) {
            System.err.println("Failed to send email to " + emailModel.getEmail() + ": " + e.getMessage());
            e.printStackTrace();
        }
    }
    
    public void sendEmail(String recipient, String subject, String message) throws Exception {
        MimeMessage msg = new MimeMessage(mailSession);
        
        // 修正原错误:设置发件人而非收件人
        msg.setFrom(new InternetAddress("xyz@gmail.com"));
        msg.addRecipient(Message.RecipientType.TO, new InternetAddress(recipient)); 
        msg.setSubject(subject);
        msg.setText(message);
        
        Transport.send(msg);
        System.out.println("Email sent successfully to: " + recipient); 
    }
}

二、关键优化点

  • 泛型修正:移除EmailBuffer错误的泛型声明,确保使用你定义的EmailModel类型。
  • 线程池复用:用线程池管理邮件发送任务,避免频繁创建销毁线程,提升性能。
  • Session复用:JavaMail的Session线程安全,仅初始化一次,减少资源消耗。
  • 异常处理优化:处理线程中断信号实现优雅停止,捕获邮件发送异常并记录错误信息。
  • 队列容量调整:LinkedBlockingQueue初始容量从1调整为10,避免生产者频繁阻塞。
  • 发件人修正:原代码错误地将收件人设为发件人,修正后确保邮件合规性。

三、使用示例(生产者类)

package com.webservices.worker;

import com.webservices.buffer.EmailBuffer;
import com.webservices.model.EmailModel;

public class EmailProducer implements Runnable {
    private EmailBuffer buffer;

    public EmailProducer(EmailBuffer buffer) {
        this.buffer = buffer;
    }

    @Override
    public void run() {
        try {
            // 模拟从API获取的验证后数据
            EmailModel email1 = new EmailModel("user1@example.com", "Welcome", "Hello User1!");
            buffer.put(email1);
            Thread.sleep(2000);
            
            EmailModel email2 = new EmailModel("user2@example.com", "Notification", "Your order is shipped.");
            buffer.put(email2);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    public static void main(String[] args) {
        EmailBuffer buffer = new EmailBuffer();
        new Thread(new EmailProducer(buffer)).start();
        new Thread(new EmailConsumer(buffer)).start();
    }
}

内容的提问来源于stack exchange,提问作者RZA

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 06:09:55