如何让Spring Cloud Stream仅在收到外部系统数据时向Kafka Topic发消息
解决方案
你当前使用的无参Supplier默认会被Spring Cloud Stream自动配置轮询调度器,默认间隔1秒触发,所以才会定时生成数据发送,这种模式本身是给定时生成数据的场景设计的,不适合你这种UDP数据包处理完成的外部事件触发场景,有两种成熟改造方案可选:
方案1:使用StreamBridge主动发送(最适配UDP监听场景,推荐)
不需要定义Supplier Bean,直接在你UDP数据包处理完成的业务逻辑里,调用StreamBridge的发送方法,即可把处理好的数据主动推送到Kafka。
代码示例
// 注入Spring Cloud Stream提供的StreamBridge工具 @Autowired private StreamBridge streamBridge; // 这是你原有处理UDP数据包的业务方法 public void afterUdpProcess(PacketDataPojo processedData) { // 原有UDP数据处理逻辑 // ... // 处理完成后主动调用发送方法 boolean sendStatus = streamBridge.send("data-out-0", processedData); log.info("数据发送结果:{},发送内容:{}", sendStatus, processedData); }
配置保持和原有一致即可
spring.cloud.function.definition=data spring.cloud.stream.bindings.data-out-0.destination=data-stream
如果不想预定义binding,也可以直接在send方法里指定目标Topic名:streamBridge.send("data-stream", processedData),对应可以删除配置里的binding配置。
方案2:使用反应式Supplier<Flux<PacketDataPojo>>实现事件触发
如果你希望依然保留Supplier的定义风格,可以返回反应式Flux,在UDP数据处理完成后往Flux推送数据,Spring Cloud Stream会自动把推送的元素发送到Kafka。
代码示例
// 定义缓冲队列,存放处理完成待发送的数据包 private final Sinks.Many<PacketDataPojo> packetSink = Sinks.many().multicast().onBackpressureBuffer(); @Bean public Supplier<Flux<PacketDataPojo>> data() { return () -> packetSink.asFlux(); } // 这是你原有处理UDP数据包的业务方法 public void afterUdpProcess(PacketDataPojo processedData) { // 原有UDP数据处理逻辑 // ... // 处理完成后往缓冲队列推送数据 packetSink.tryEmitNext(processedData); }
配置完全不需要修改,和你原有配置一致即可
内容的提问来源于stack exchange,提问作者Dnyaneshwar
相关产品推荐
相关产品推荐

