【问题标题】:Spring cloud function Cosuming rabbitmq queue - Dispatcher has no subscribers for channel春季云功能使用rabbitmq队列 - 调度程序没有频道订阅者
【发布时间】:2020-04-15 05:44:36
【问题描述】:

给定:

  cloud:
    stream:
      rabbit:
      bindings:
        inboundApolloLookupVehicleChannel:
          destination: fed.apollo-vehicle-lookup-test
          group: apollo-mngt-group
          consumer:
            missingQueuesFatal: true
            prefetch: 25
            autoBindDlq: true
            maxAttempts: 1
            republishToDlq: true
            requeueRejected: false
            durableSubscription: true
            maxConcurrency: 6              
      function:
        definition: consume  

和:

@SpringBootApplication
@EnableEurekaClient
@EnableBinding(LookupMessageChannel.class)
public class ApolloLookupServiceApplication {

    public static void main(String[] args) {
        SpringApplication.run(ApolloLookupServiceApplication.class, args);
    }
}


public interface LookupMessageChannel {

    String INBOUND_LOOKUP_VEHICLE_CHANNEL = "inboundApolloLookupVehicleChannel";

    @Input(INBOUND_LOOKUP_VEHICLE_CHANNEL)
    SubscribableChannel inboundApolloLookupVehicleChannel();

}

@Service
@MessageEndpoint
@RequiredArgsConstructor
public class ApolloVehicleLookupService {

    private final ApolloVehicleLookUpRepository apolloVehicleLookUpRepository;
    private static final Logger LOGGER = LoggerFactory.getLogger(ApolloVehicleLookupService.class);

    @Bean
    public Consumer<Flux<ApolloVehicleLookUp>> consume() {
        return stream -> stream             
                .flatMap(this.apolloVehicleLookUpRepository::save)
                .subscribe(value -> {
                        LOGGER.info("stored value: " + value.toString());
                });
    }

}

POM:



    <parent>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-parent</artifactId>
        <version>2.2.6.RELEASE</version>
        <relativePath />
        <!-- lookup parent from repository -->
    </parent>

    <properties>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
        <project.reporting.outputEncoding>UTF-8</project.reporting.outputEncoding>
        <java.version>14</java.version>
        <spring-cloud.version>Hoxton.SR3</spring-cloud.version>
        <os-maven-plugin.version>1.5.0.Final</os-maven-plugin.version>
        <spotify-docker-maven.version>1.2.0</spotify-docker-maven.version>      
        <os.detected.classifier>linux-x86_64</os.detected.classifier>
    </properties>

  <dependencies>

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-data-mongodb-reactive</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-rsocket</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-webflux</artifactId>
        </dependency>
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <optional>true</optional>
        </dependency>
        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-actuator</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-bus-amqp</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-config</artifactId>
        </dependency>
        <dependency>
            <groupId>org.springframework.cloud</groupId>
            <artifactId>spring-cloud-starter-netflix-eureka-client</artifactId>
        </dependency>
  </dependencies>


为什么会出现以下异常?

[payload=org.springframework.messaging.MessageDeliveryException: Dispatcher 没有频道“apollo-lookup-service-1.inboundApolloLookupVehicleChannel”的订阅者。;嵌套异常是 org.springframework.integration.MessageDispatchingException: Dispatcher 没有订阅者,failedMessage=GenericMessage [payload=byte[439]

我只是想使用该队列中的消息。

任何帮助将不胜感激,谢谢。

【问题讨论】:

    标签: java spring spring-cloud-stream spring-cloud-function


    【解决方案1】:

    异常简单意味着您在LookupMessageChannel 中定义的频道没有订阅者,为什么要这样,因为您没有定义任何订阅者。在我看来,您正在尝试使用功能方法,而您的流应用程序的其余部分使用几乎已弃用的旧配置。

    不要误会我的意思,但目前您的应用几乎没有什么问题,因此很难确定您到底要做什么,所以。 . .

    请考虑关注this quick start(顶部 5 分钟),让您了解正确的功能方法,然后随时跟进其他问题。

    【讨论】:

    • 感谢@Oleg,我们正在迁移服务,但很多地方都出错了。现在它更整洁了。我按照快速入门,它开始工作。这是我的新问题:如果我为旧服务生成了一条旧消息,并且我将它的有效负载复制并粘贴到新队列中,它会很好地处理它,但是如果我将带有标头的完整消息移动到新队列我' m 得到以下异常: 原因:org.springframework.integration.MessageDispatchingException:Dispatcher 没有订阅者。 这是因为标头吗?我应该使用生产者函数生成消息吗?
    猜你喜欢
    • 2011-12-18
    • 1970-01-01
    • 2015-11-02
    • 1970-01-01
    • 2018-02-19
    • 2020-04-22
    • 1970-01-01
    • 2015-05-18
    • 1970-01-01
    相关资源
    最近更新 更多