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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:02:43