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

如何用IntegrationFlow实现Spring Cloud Stream Kafka消息发送并解决订阅者缺失异常

问题描述

原基于Spring Cloud Stream Kafka的Supplier代码可正常运行:

@SpringBootApplication
public class DemoApplication {

    public static void main(String[] args) {
        SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic");
    }
    @Bean
    Supplier<String> outbound() {
        return () -> {
            return LocalTime.now().toString();
        };
    }
}

尝试改用IntegrationFlow实现相同功能并添加转换器时,抛出异常MessageDispatchingException: Dispatcher has no subscribers,用户代码如下:

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.Pollers;
import java.time.LocalTime;

@SpringBootApplication
public class DemoApplication {

    public static void main(String[] args) {
        SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic");
    }
    @Bean
    IntegrationFlow myFlow() {
        return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000)))
                .transform(m -> {
                    System.out.println("my transformer");
                    return m;
                })
                .channel("outbound")
                .get();
    }

    String myPoller() {
            return LocalTime.now() + " value";
    }
}

Gradle配置:

plugins {
    id 'java'
    id 'org.springframework.boot' version '3.2.9'
    id 'io.spring.dependency-management' version '1.1.6'
}

group = 'com.example'
version = '0.0.1-SNAPSHOT'

java {
    toolchain {
        languageVersion = JavaLanguageVersion.of(21)
    }
}

repositories {
    mavenCentral()
}
	ext {
    set('springCloudVersion', "2023.0.3")
}

dependencies {
    implementation 'org.springframework.cloud:spring-cloud-stream'
    implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka'
    implementation 'org.springframework.kafka:spring-kafka'
    testImplementation 'org.springframework.boot:spring-boot-starter-test'
    testImplementation 'org.springframework.cloud:spring-cloud-stream-test-binder'
    testImplementation 'org.springframework.kafka:spring-kafka-test'
    testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}

dependencyManagement {
    imports {
        mavenBom "org.springframework.cloud:spring-cloud-dependencies:${springCloudVersion}"
    }
}

tasks.named('test') {
    useJUnitPlatform()
}
错误原因

代码中.channel("outbound")指定的是普通Spring Integration通道,该通道未与Kafka出站绑定关联,因此没有订阅者,导致消息无法分发。而原Supplier方式是通过Spring Cloud Stream的函数式绑定自动创建了与Kafka关联的出站通道(outbound-out-0),消息能正常发送到Kafka。

解决方案

以下两种方式均可实现需求,将IntegrationFlow的输出连接到Spring Cloud Stream的Kafka出站绑定:

方式1:直接绑定到自动创建的出站通道

利用配置中指定的outbound-out-0通道名(对应spring.cloud.stream.bindings.outbound-out-0.destination=out-topic),将IntegrationFlow的输出指向该通道:

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.Pollers;
import java.time.LocalTime;

@SpringBootApplication
public class DemoApplication {

    public static void main(String[] args) {
        SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic");
    }

    @Bean
    IntegrationFlow myFlow() {
        return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000)))
                .transform(m -> {
                    System.out.println("my transformer");
                    return m;
                })
                .channel("outbound-out-0") // 替换为Spring Cloud Stream创建的出站通道
                .get();
    }

    String myPoller() {
        return LocalTime.now() + " value";
    }
}

方式2:使用StreamBridge发送(更灵活)

通过StreamBridge直接发送消息到指定绑定目标,无需依赖固定通道名:

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.stream.function.StreamBridge;
import org.springframework.context.annotation.Bean;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.Pollers;
import java.time.LocalTime;

@SpringBootApplication
public class DemoApplication {

    public static void main(String[] args) {
        SpringApplication.run(DemoApplication.class, "--spring.cloud.stream.bindings.outbound-out-0.destination=out-topic");
    }

    @Bean
    IntegrationFlow myFlow(StreamBridge streamBridge) {
        return IntegrationFlow.fromSupplier(this::myPoller, p -> p.poller(Pollers.fixedDelay(5000)))
                .transform(m -> {
                    System.out.println("my transformer");
                    return m;
                })
                .handle(message -> streamBridge.send("outbound-out-0", message.getPayload())) // 通过StreamBridge发送消息
                .get();
    }

    String myPoller() {
        return LocalTime.now() + " value";
    }
}
说明
  • 方式1中,outbound-out-0是Spring Cloud Stream根据绑定配置自动创建的出站通道,已与Kafka绑定器关联,消息会被发送到配置的out-topic。
  • 方式2的StreamBridge支持动态指定发送目标,适合需要根据业务逻辑切换主题的场景,同样会利用Spring Cloud Stream的绑定配置完成Kafka消息发送。

内容的提问来源于stack exchange,提问作者mike27

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 06:07:04