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
相关产品推荐
相关产品推荐

