Camel处理中JVM崩溃后如何从失败点恢复消息处理?
解决方案:Camel路由拆分任务的断点续传实现
核心问题分析
你遇到的情况是JVM异常崩溃而非正常停机,所以DefaultShutdownStrategy完全起不到作用——这个策略只在应用收到正常停止信号时触发,用于等待 inflight 消息处理完成。要实现断点续传,必须通过持久化处理进度或重复请求过滤来避免重复处理已完成的部分。
方案一:基于进度持久化的断点续传(精准控制)
这个方案直接针对你的需求,能从崩溃时的位置继续处理,步骤如下:
修改File组件配置,取消自动删除
把原路由中delete=true去掉,改用处理完成后手动移动文件:<from uri="file://{{env:TEST}}/FOLDER?exclude=.*filepart&delay={{delay}}&readLock=changed&readLockCheckInterval={{readLockCheckInterval}}&readLockMinAge={{readLockMinAge}}&initialDelay={{initialDelay}}"/>添加进度检查与记录逻辑
新增两个Bean:progressChecker(读取上次崩溃的进度)和progressRecorder(记录当前处理进度),修改路由如下:<route id="XYZ"> <from uri="file://{{env:TEST}}/FOLDER?exclude=.*filepart&delay={{delay}}&readLock=changed&readLockCheckInterval={{readLockCheckInterval}}&readLockMinAge={{readLockMinAge}}&initialDelay={{initialDelay}}"/> <!-- 检查是否有未完成的进度,返回起始索引(比如501) --> <to uri="bean:progressChecker" method="getStartIndex"/> <!-- 从起始索引开始拆分数据 --> <split> <method ref="MethodA" method="splitFromIndex(${body}, ${header.startIndex})"/> <!-- 每条消息发送后记录当前索引 --> <setHeader headerName="currentIndex"> <simple>${exchangeProperty.CamelSplitIndex}</simple> </setHeader> <to uri="activemq:testQueue"/> <to uri="bean:progressRecorder" method="updateIndex(${header.currentIndex})"/> </split> <!-- 全部处理完成后清除进度标记,并移动原文件到归档目录 --> <to uri="bean:progressRecorder" method="clearProgress"/> <to uri="bean:fileArchiver" method="moveToProcessed"/> </route>progressChecker:读取本地文件/数据库中的进度记录,若存在则返回上次中断的索引+1(比如501),否则返回0。MethodA.splitFromIndex:修改拆分方法,支持从指定索引开始跳过已处理的条目,只返回未处理的数据。progressRecorder:每处理一条就将当前索引写入持久化存储(比如本地JSON文件、数据库表),崩溃重启后可读取该值继续。
方案二:幂等消费者过滤重复消息(无需修改拆分逻辑)
如果不想改动拆分方法,可以用Camel的幂等消费者功能,自动跳过已发送到MQ的消息:
配置幂等仓库
用数据库或LevelDB存储已处理的消息唯一标识,示例用JDBC仓库:<bean id="idempotentRepo" class="org.apache.camel.processor.idempotent.jdbc.JdbcMessageIdRepository"> <constructor-arg ref="dataSource"/> <constructor-arg value="processed_msg_ids"/> <!-- 数据库表名,需提前创建 --> </bean>修改路由添加幂等校验
假设每条拆分后的数据有唯一ID字段(比如dataId),路由配置如下:<route id="XYZ"> <from uri="file://{{env:TEST}}/FOLDER?exclude=.*filepart&delay={{delay}}&readLock=changed&readLockCheckInterval={{readLockCheckInterval}}&readLockMinAge={{readLockMinAge}}&initialDelay={{initialDelay}}"/> <to uri="bean:testprocess"/> <split> <method ref="MethodA" method="split"/> <!-- 设置消息唯一标识 --> <setHeader headerName="msgUniqueId"> <simple>${body.dataId}</simple> </setHeader> <!-- 幂等校验,已处理的消息直接跳过 --> <idempotentConsumer messageIdRepositoryRef="idempotentRepo" skipDuplicate="true"> <header>msgUniqueId</header> <to uri="activemq:testQueue"/> </idempotentConsumer> </split> <to uri="bean:fileArchiver" method="moveToProcessed"/> </route>重启后,拆分出的消息会先检查幂等仓库,已存在的(前500条)直接跳过,只处理未记录的部分。
方案三:事务回滚(避免部分发送,非断点续传)
如果你的场景不允许出现部分消息发送的情况,可以用ActiveMQ事务确保要么全部发送成功,要么全部回滚,但重启后会重新处理全部数据:
配置JMS事务管理器
<bean id="jmsTxManager" class="org.springframework.jms.connection.JmsTransactionManager"> <property name="connectionFactory" ref="activemqConnectionFactory"/> </bean>路由开启事务
<route id="XYZ" transacted="jmsTxManager"> <from uri="file://{{env:TEST}}/FOLDER?exclude=.*filepart&delay={{delay}}&delete=false&readLock=changed&readLockCheckInterval={{readLockCheckInterval}}&readLockMinAge={{readLockMinAge}}&initialDelay={{initialDelay}}"/> <to uri="bean:testprocess"/> <split streaming="true"> <method ref="MethodA" method="split"/> <to uri="activemq:testQueue"/> </split> <to uri="bean:fileArchiver" method="moveToProcessed"/> </route>JVM崩溃时,未提交的事务会回滚,已发送的500条消息会被ActiveMQ收回,重启后重新处理全部2000条。
内容的提问来源于stack exchange,提问作者Niveditha N
相关产品推荐
相关产品推荐

