Apache Camel如何将失败消息持久化到PostgreSQL死信表?
Apache Camel 持久化失败消息到PostgreSQL的配置问题
我用Apache Camel搭建集成系统,需要持久化失败消息,方便后续通过UI重新投递处理。目前已经实现了把失败消息发送到Kafka死信队列的逻辑,但在持久化到PostgreSQL时遇到了配置问题。
Kafka死信队列示例配置
- routeConfiguration: errorHandler: id: errorHandler-d669 deadLetterChannel: deadLetterUri: kafka:failed-messages?brokers=kafka:9092 redeliveryPolicy: maximumRedeliveries: 2 redeliveryDelay: '1000' backOffMultiplier: 2 useExponentialBackOff: true retryAttemptedLogLevel: WARN id: redeliveryPolicy-215f level: ERROR id: deadLetterChannel-6c51 id: deadletter-kafka
该配置会将完整Exchange消息发送到名为failed-messages的Kafka主题,生成的队列条目示例如下:
Exchange[Id: 7384875AE37D9B7-0000000000000004, ExchangePattern: InOnly, Properties: {CamelMessageHistory=[DefaultMessageHistory[routeId=route-5893, node=to-d19e]], CamelToEndpoint=log://deadletter-logger?level=WARN&showAll=true}, Headers: {CamelMessageTimestamp=1703585126556, Custom-Header-1=[B@1509b733, Custom-Header-2=[B@d9561f8, kafka.HEADERS=RecordHeaders(headers = [RecordHeader(key = Custom-Header-1, value = [79, 110, 101]), RecordHeader(key = Custom-Header-2, value = [84, 119, 111])], isReadOnly = false), kafka.OFFSET=11815, kafka.PARTITION=0, kafka.TIMESTAMP=1703585126556, kafka.TOPIC=failed-messages}, BodyType: String, Body: It's me, Mario, on 2023-12-26T10:05:23 - Enriched]
PostgreSQL持久化的配置问题
我尝试配置死信通道将消息存入PostgreSQL,但不知道该如何指定读取失败消息的表达式,当前配置如下:
- routeConfiguration: errorHandler: id: errorHandler-d679 deadLetterChannel: deadLetterUri: >- sql:insert into camel_deadletter (exchange_message) values (:#${SOMETHING}) redeliveryPolicy: maximumRedeliveries: 2 redeliveryDelay: '1000' backOffMultiplier: 2 useExponentialBackOff: true retryAttemptedLogLevel: WARN id: redeliveryPolicy-216f level: ERROR id: deadLetterChannel-6b51 id: deadletter-sql
需要替换上述配置中的SOMETHING,但不确定该填写什么内容。
解决方案
要完整持久化Exchange消息并支持后续重新投递,需要将整个Exchange对象序列化后存入数据库,以下是两种可行配置方式:
方式1:直接在SQL URI中序列化Exchange
利用Camel表达式结合Jackson直接将Exchange转为JSON字符串,填入SQL语句中:
deadLetterUri: >- sql:insert into camel_deadletter (exchange_message) values (:#${exchange.usingPattern(com.fasterxml.jackson.databind.ObjectMapper).writeValueAsString(exchange)})
注意:项目中需提前引入Jackson相关依赖,否则序列化会失败。
方式2:通过Bean处理器序列化(更易维护)
先定义一个Bean处理器负责序列化Exchange,再在死信通道中调用该处理器,最后将序列化后的内容存入数据库:
1. 修改死信通道配置
- routeConfiguration: errorHandler: id: errorHandler-d679 deadLetterChannel: deadLetterUri: sql:insert into camel_deadletter (exchange_message) values (:#body) onPrepareFailure: processor: bean: ref: exchangeSerializerBean method: serializeExchange redeliveryPolicy: maximumRedeliveries: 2 redeliveryDelay: '1000' backOffMultiplier: 2 useExponentialBackOff: true retryAttemptedLogLevel: WARN id: redeliveryPolicy-216f level: ERROR id: deadLetterChannel-6b51 id: deadletter-sql
2. 实现序列化Bean
import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.camel.Exchange; import org.springframework.stereotype.Component; @Component("exchangeSerializerBean") public class ExchangeSerializer { private final ObjectMapper objectMapper = new ObjectMapper(); public String serializeExchange(Exchange exchange) throws Exception { return objectMapper.writeValueAsString(exchange); } }
额外建议
- PostgreSQL表的
exchange_message字段建议使用JSONB类型(比TEXT更适合存储JSON,支持查询操作),若仅需存储字符串则用TEXT类型。 - 如果不需要完整Exchange,仅需消息体和头部,可分别引用
body和headers:
deadLetterUri: >- sql:insert into camel_deadletter (message_body, message_headers) values (:#${body}, :#${headers})
但这种方式无法完整恢复Exchange用于重新投递,若需后续重新处理,建议存储完整序列化的Exchange。
内容的提问来源于stack exchange,提问作者Damir Palinić
相关产品推荐
相关产品推荐

