如何用Spring Integration实现Feed入站通道适配器对接MongoDB与Kafka
嘿,这个需求用Spring Integration的Java DSL来实现简直是量身定做!你提到的**发布/订阅通道(PublishSubscribeChannel)**正是核心——它能让拉取到的RSS Feed消息同时流向MongoDB存储和Kafka发送两个分支,完美覆盖故障处理、审计以及后续消费的场景。
- 用
Feed.inboundAdapter定时拉取指定RSS Feed的内容; - 将拉取到的消息发送到发布/订阅通道,实现消息的广播分发;
- 两个独立的处理分支:
- 分支1:将RSS条目转换为自定义实体类,存储到MongoDB数据库(用于故障回溯、审计);
- 分支2:将同一条RSS条目序列化为JSON格式,发送到指定Kafka主题供下游服务消费。
首先需要确保项目中引入了必要的依赖:spring-integration-rss、spring-integration-mongodb、spring-integration-kafka、spring-boot-starter-data-mongodb以及Kafka相关的Spring Boot starter。
1. 自定义RSS实体类(对应MongoDB文档)
import org.springframework.data.annotation.Id; import org.springframework.data.mongodb.core.mapping.Document; import java.util.Date; @Document(collection = "rss_feed_entries") public class RssFeedEntry { @Id private String id; private String title; private String link; private Date publishDate; private String description; // 建议用Lombok的@Data注解自动生成getter、setter和构造函数 // 这里为了兼容性手动省略,实际项目中可以简化 public String getId() { return id; } public void setId(String id) { this.id = id; } public String getTitle() { return title; } public void setTitle(String title) { this.title = title; } public String getLink() { return link; } public void setLink(String link) { this.link = link; } public Date getPublishDate() { return publishDate; } public void setPublishDate(Date publishDate) { this.publishDate = publishDate; } public String getDescription() { return description; } public void setDescription(String description) { this.description = description; } }
2. Spring Integration 核心配置类
import com.rometools.rome.feed.synd.SyndEntry; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.integration.annotation.EnableIntegration; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.dsl.Pollers; import org.springframework.integration.feed.inbound.FeedEntryMessageSource; import org.springframework.integration.transformer.Transformers; import org.springframework.messaging.MessageChannel; import org.springframework.kafka.core.KafkaTemplate; import java.net.URL; @Configuration @EnableIntegration public class RssFeedIntegrationConfig { private final MongoTemplate mongoTemplate; private final KafkaTemplate<String, String> kafkaTemplate; @Value("${rss.feed.url}") private String rssFeedUrl; @Value("${kafka.topic.rss-feed}") private String rssKafkaTopic; // 构造函数注入依赖 public RssFeedIntegrationConfig(MongoTemplate mongoTemplate, KafkaTemplate<String, String> kafkaTemplate) { this.mongoTemplate = mongoTemplate; this.kafkaTemplate = kafkaTemplate; } // 定义发布/订阅通道,实现消息广播 @Bean public MessageChannel rssFeedPublishSubscribeChannel() { return new org.springframework.integration.channel.PublishSubscribeChannel(); } // RSS Feed入站流:定时拉取Feed数据并发送到发布订阅通道 @Bean public IntegrationFlow rssFeedInboundFlow() throws Exception { FeedEntryMessageSource feedSource = new FeedEntryMessageSource(new URL(rssFeedUrl), "rss-feed-poller-marker"); return IntegrationFlows.from(feedSource, spec -> spec.poller(Pollers.fixedDelay(60000) // 每分钟拉取一次,可配置 .maxMessagesPerPoll(10))) // 每次拉取最多10条消息 .channel(rssFeedPublishSubscribeChannel()) .get(); } // 分支1:将Feed数据存储到MongoDB @Bean public IntegrationFlow rssFeedToMongoDbFlow() { return IntegrationFlows.from(rssFeedPublishSubscribeChannel()) .transform(this::convertSyndEntryToRssFeedEntry) // 转换为自定义实体类 .handle(MongoDb.outboundAdapter(mongoTemplate) .collectionName("rss_feed_entries")) // 指定MongoDB集合 .get(); } // 分支2:将Feed数据序列化为JSON并发送到Kafka @Bean public IntegrationFlow rssFeedToKafkaFlow() { return IntegrationFlows.from(rssFeedPublishSubscribeChannel()) .transform(this::convertSyndEntryToRssFeedEntry) .transform(Transformers.toJson()) // 转换为JSON字符串 .handle(Kafka.outboundChannelAdapter(kafkaTemplate) .topic(rssKafkaTopic)) // 指定Kafka主题 .get(); } // 自定义转换逻辑:把Rome库的SyndEntry转为自定义实体类 private RssFeedEntry convertSyndEntryToRssFeedEntry(SyndEntry syndEntry) { RssFeedEntry entry = new RssFeedEntry(); entry.setId(syndEntry.getUri()); entry.setTitle(syndEntry.getTitle()); entry.setLink(syndEntry.getLink()); entry.setPublishDate(syndEntry.getPublishedDate()); entry.setDescription(syndEntry.getDescription() != null ? syndEntry.getDescription().getValue() : null); return entry; } }
发布/订阅通道的必要性:
完全需要!这个通道会把同一份RSS消息广播给所有订阅的处理器,MongoDB存储和Kafka发送两个分支会各自收到完整的消息副本,互不干扰,正好符合你同时存储和转发的需求。RSS拉取标记的持久化:
默认情况下,FeedEntryMessageSource的拉取标记存在内存中,重启服务后会重复拉取旧消息。如果需要持久化标记,可以添加MongoDB的MetadataStore:@Bean public org.springframework.integration.metadata.MetadataStore metadataStore(MongoTemplate mongoTemplate) { return new org.springframework.integration.mongodb.metadata.MongoDbMetadataStore(mongoTemplate); }然后在
rssFeedInboundFlow的FeedEntryMessageSource中指定:feedSource.setMetadataStore(metadataStore());故障处理与审计:
可以为每个分支添加错误处理逻辑,比如配置重试机制或死信队列:// 示例:为MongoDB分支添加重试建议 .handle(MongoDb.outboundAdapter(mongoTemplate).collectionName("rss_feed_entries"), e -> e.advice(retryAdvice()))失败的消息可以存入MongoDB的死信集合或Kafka的死信主题,方便后续审计和重试。
在application.properties中添加MongoDB和Kafka的连接配置:
# MongoDB配置 spring.data.mongodb.uri=mongodb://localhost:27017/rss_feed_db # Kafka配置 spring.kafka.bootstrap-servers=localhost:9092 spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer=org.apache.kafka.common.serialization.StringSerializer # RSS和Kafka主题配置 rss.feed.url=https://example.com/your-rss-feed.xml kafka.topic.rss-feed=rss-feed-processing-topic
内容的提问来源于stack exchange,提问作者Swordfish

