Spring Integration:IMAP邮件统计重排器分组释放问题求解
解决Spring Integration邮件统计中的Resequencer分组与释放问题
我来帮你梳理下这个需求的解决方案,针对你用Spring Integration实现每周邮件同步并统计主题数量的场景,咱们一步步来解决核心问题:
1. 给邮件消息添加Correlation ID和Sequence Number
Resequencer依赖这两个标识来分组和管理消息,而MimeMessage本身没有这些属性,所以我们需要通过Transformer给每个消息注入这两个头信息:
- Correlation ID:因为是每周一次同步,所有本次同步的邮件属于同一个分组,我们可以用每周唯一的标识(比如当前周的年份+周数,如
2024-W23),确保同周的消息都分到一组。 - Sequence Number:用递增的序号即可,保证每个消息的序号唯一,不需要严格排序的话,甚至可以用随机数,但递增更稳妥。
写一个简单的Transformer实现:
import org.springframework.integration.transformer.GenericTransformer; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import javax.mail.internet.MimeMessage; import java.util.concurrent.atomic.AtomicInteger; public class MailHeaderTransformer implements GenericTransformer<Message<MimeMessage>, Message<MimeMessage>> { private final String weeklyCorrelationId; private final AtomicInteger sequenceCounter = new AtomicInteger(1); public MailHeaderTransformer(String weeklyCorrelationId) { this.weeklyCorrelationId = weeklyCorrelationId; } @Override public Message<MimeMessage> transform(Message<MimeMessage> message) { return MessageBuilder.fromMessage(message) .setHeader("correlationId", weeklyCorrelationId) .setHeader("sequenceNumber", sequenceCounter.getAndIncrement()) .build(); } }
然后在配置里添加这个Transformer,放在邮件适配器和Resequencer之间:
<int:transformer input-channel="receiveChannel" output-channel="transformedMailChannel" ref="mailHeaderTransformer"/> <bean id="mailHeaderTransformer" class="com.yourpackage.MailHeaderTransformer"> <!-- 用SpEL生成每周唯一的Correlation ID --> <constructor-arg value="#{T(java.time.LocalDate).now().getWeekYear() + '-W' + T(java.time.LocalDate).now().get(IsoFields.WEEK_OF_WEEK_BASED_YEAR)}"/> </bean> <int:channel id="transformedMailChannel"/>
2. 配置Resequencer的释放策略(或改用Aggregator更合适)
你的核心需求是收集完所有本次同步的邮件后再生成统计,其实用Aggregator比Resequencer更贴合场景——Resequencer侧重乱序消息的排序输出,而Aggregator专门用于收集一组消息后输出汇总结果(比如邮件列表)。
方案:用Aggregator收集所有邮件,同步完成后触发释放
我们可以监听邮件轮询的结束事件,发送一个释放信号给Aggregator,让它把收集到的所有邮件一次性输出:
第一步:配置Aggregator
<int:aggregator input-channel="transformedMailChannel" output-channel="statisticsChannel" correlation-strategy-expression="headers.correlationId" <!-- 释放条件:收到带有releaseSignal头的消息 --> release-strategy-expression="headers.containsKey('releaseSignal')" message-store="inMemoryMessageStore"> <!-- 聚合策略:直接返回所有消息的列表 --> <int:aggregation-strategy expression="payload"/> </int:aggregator> <bean id="inMemoryMessageStore" class="org.springframework.integration.store.SimpleMessageStore"/> <int:channel id="statisticsChannel"/>
第二步:监听轮询结束事件,发送释放信号
写一个监听器,监听邮件适配器的轮询终止事件,事件触发后发送释放信号:
import org.springframework.context.ApplicationListener; import org.springframework.integration.endpoint.SourcePollingChannelAdapter; import org.springframework.integration.event.inbound.PollTerminatedEvent; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.MessageBuilder; public class MailSyncCompletionListener implements ApplicationListener<PollTerminatedEvent> { private final MessageChannel transformedMailChannel; private final String weeklyCorrelationId; public MailSyncCompletionListener(MessageChannel transformedMailChannel, String weeklyCorrelationId) { this.transformedMailChannel = transformedMailChannel; this.weeklyCorrelationId = weeklyCorrelationId; } @Override public void onApplicationEvent(PollTerminatedEvent event) { SourcePollingChannelAdapter adapter = (SourcePollingChannelAdapter) event.getSource(); // 确认是我们的邮件适配器触发的轮询结束 if ("mailClient".equals(adapter.getComponentName())) { Message<String> releaseSignal = MessageBuilder.withPayload("") .setHeader("correlationId", weeklyCorrelationId) .setHeader("releaseSignal", true) .build(); transformedMailChannel.send(releaseSignal); } } }
在配置里注册这个监听器:
<bean id="mailSyncCompletionListener" class="com.yourpackage.MailSyncCompletionListener"> <constructor-arg ref="transformedMailChannel"/> <constructor-arg value="#{T(java.time.LocalDate).now().getWeekYear() + '-W' + T(java.time.LocalDate).now().get(IsoFields.WEEK_OF_WEEK_BASED_YEAR)}"/> </bean>
第三步:添加统计服务
最后,在statisticsChannel后面加一个Service Activator,处理汇总后的邮件列表,统计相同主题的数量:
import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.handler.annotation.Payload; import javax.mail.internet.MimeMessage; import java.util.List; import java.util.Map; import java.util.stream.Collectors; public class MailStatisticsService { @ServiceActivator(inputChannel = "statisticsChannel") public void generateAndSendStatistics(@Payload List<MimeMessage> messages) throws Exception { // 统计相同主题的邮件数量 Map<String, Long> subjectCountMap = messages.stream() .map(mimeMessage -> { try { return mimeMessage.getSubject(); } catch (Exception e) { return "未知主题"; } }) .collect(Collectors.groupingBy(subject -> subject, Collectors.counting())); // 这里可以把统计结果通过邮件、通知等方式发送出去 subjectCountMap.forEach((subject, count) -> System.out.printf("主题「%s」共收到 %d 封邮件%n", subject, count) ); } }
配置这个服务:
<bean id="mailStatisticsService" class="com.yourpackage.MailStatisticsService"/> <int:service-activator input-channel="statisticsChannel" ref="mailStatisticsService"/>
整合后的完整配置
把所有部分整合起来,你的integration-context.xml大概是这样:
<int:channel id="startMailSync"/> <int:control-bus id="start" input-channel="startMailSync"/> <int:channel id="receiveChannel" datatype="javax.mail.internet.MimeMessage"/> <int:channel id="transformedMailChannel"/> <int:channel id="statisticsChannel"/> <int-mail:inbound-channel-adapter id="mailClient" channel="receiveChannel" java-mail-properties="javaMailProperties" store-uri="imaps://[user]:[password]@mail.it/INBOX" should-mark-messages-as-read="true" should-delete-messages="false" mail-filter-expression="from[0].address matches 'sender@sender.it'" auto-startup="false"> <int:poller trigger="runOnceTrigger" max-messages-per-poll="100"/> </int-mail:inbound-channel-adapter> <util:properties id="javaMailProperties"> <prop key="mail.imap.socketFactory.class">javax.net.ssl.SSLSocketFactory</prop> <prop key="mail.imap.socketFactory.fallback">false</prop> <prop key="mail.store.protocol">imaps</prop> <prop key="mail.debug">false</prop> </util:properties> <bean id="runOnceTrigger" class="org.springframework.scheduling.support.PeriodicTrigger"> <constructor-arg value="0"/> <property name="initialDelay" value="0"/> <property name="fixedRate" value="#{T(java.lang.Long).MAX_VALUE}"/> </bean> <int:transformer input-channel="receiveChannel" output-channel="transformedMailChannel" ref="mailHeaderTransformer"/> <bean id="mailHeaderTransformer" class="com.yourpackage.MailHeaderTransformer"> <constructor-arg value="#{T(java.time.LocalDate).now().getWeekYear() + '-W' + T(java.time.LocalDate).now().get(IsoFields.WEEK_OF_WEEK_BASED_YEAR)}"/> </bean> <int:aggregator input-channel="transformedMailChannel" output-channel="statisticsChannel" correlation-strategy-expression="headers.correlationId" release-strategy-expression="headers.containsKey('releaseSignal')" message-store="inMemoryMessageStore"> <int:aggregation-strategy expression="payload"/> </int:aggregator> <bean id="inMemoryMessageStore" class="org.springframework.integration.store.SimpleMessageStore"/> <bean id="mailSyncCompletionListener" class="com.yourpackage.MailSyncCompletionListener"> <constructor-arg ref="transformedMailChannel"/> <constructor-arg value="#{T(java.time.LocalDate).now().getWeekYear() + '-W' + T(java.time.LocalDate).now().get(IsoFields.WEEK_OF_WEEK_BASED_YEAR)}"/> </bean> <bean id="mailStatisticsService" class="com.yourpackage.MailStatisticsService"/> <int:service-activator input-channel="statisticsChannel" ref="mailStatisticsService"/>
额外提示
- 如果邮件数量很大,担心内存问题,可以把
SimpleMessageStore换成JdbcMessageStore,用数据库存储消息组。 - 如果你坚持要用Resequencer,释放策略可以设置为
sequenceNumber == headers.totalMessageCount,但需要提前获取本次同步的邮件总数并注入到每个消息的头中,这种方式比Aggregator繁琐不少,更推荐用Aggregator完成收集汇总的需求。
内容的提问来源于stack exchange,提问作者Stefano Veloccia
相关产品推荐
相关产品推荐

