spring-cloud / spring-cloud/spring-cloud-stream

spring.cloud.stream.sendto.destination not honoring explicit binder environment

Open
#1,909 14 comments 0 reactions 1 assignee View on GitHub

@olegz is already working on this.

Since Feb 12, 2020.

enhancement
Dominant language
Java
Stars
1.1k
Forks
646
Avg merge
2d 3h
Merged PRs (30d)
8

Description

spring.cloud.stream.sendto.destination is not honoring the explicit binder environment(child context), and is trying to fetch default binder assignment from parent context which is resulting in following exception.

2020-02-11 16:10:07.284 ERROR 59803 --- [oundedElastic-1] o.s.integration.handler.LoggingHandler   : org.springframework.messaging.MessageHandlingException: error occurred in message handler [bean 'supplier_integrationflow.router#0' for component 'supplier_integrationflow.org.springframework.integration.config.ConsumerEndpointFactoryBean#0']; nested exception is java.lang.IllegalStateException: A default binder has been requested, but there is more than one binder available for 'org.springframework.cloud.stream.messaging.DirectWithAttributesChannel' : kafka2,kafka1, and no default binder has been set., failedMessage=GenericMessage [payload=byte[8], headers={id=8b7de95a-6f07-8ad5-b0f5-780297eed494, spring.cloud.stream.sendto.destination=topic-out, contentType=application/json, timestamp=1581466207264}]
        at org.springframework.integration.support.utils.IntegrationUtils.wrapInHandlingExceptionIfNecessary(IntegrationUtils.java:191)
        at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:187)
        at org.springframework.integration.handler.AbstractMessageHandler.onNext(AbstractMessageHandler.java:219)
        at org.springframework.integration.handler.AbstractMessageHandler.onNext(AbstractMessageHandler.java:57)
        at org.springframework.integration.endpoint.ReactiveStreamsConsumer$DelegatingSubscriber.hookOnNext(ReactiveStreamsConsumer.java:165)
        at org.springframework.integration.endpoint.ReactiveStreamsConsumer$DelegatingSubscriber.hookOnNext(ReactiveStreamsConsumer.java:148)
        at reactor.core.publisher.BaseSubscriber.onNext(BaseSubscriber.java:160)
        at reactor.core.publisher.FluxDoFinally$DoFinallySubscriber.onNext(FluxDoFinally.java:123)
        at reactor.core.publisher.EmitterProcessor.drain(EmitterProcessor.java:426)
        at reactor.core.publisher.EmitterProcessor.onNext(EmitterProcessor.java:268)
        at reactor.core.publisher.FluxCreate$BufferAsyncSink.drain(FluxCreate.java:793)
        at reactor.core.publisher.FluxCreate$BufferAsyncSink.next(FluxCreate.java:718)
        at reactor.core.publisher.FluxCreate$SerializedSink.next(FluxCreate.java:153)
        at org.springframework.integration.channel.FluxMessageChannel.doSend(FluxMessageChannel.java:63)
        at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:453)
        at org.springframework.integration.channel.AbstractMessageChannel.send(AbstractMessageChannel.java:403)
        at org.springframework.integration.channel.FluxMessageChannel.lambda$subscribeTo$2(FluxMessageChannel.java:83)
        at reactor.core.publisher.FluxPeekFuseable$PeekFuseableSubscriber.onNext(FluxPeekFuseable.java:189)
        at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.runAsync(FluxPublishOn.java:398)
        at reactor.core.publisher.FluxPublishOn$PublishOnSubscriber.run(FluxPublishOn.java:484)
        at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:84)
        at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:37)
        at java.util.concurrent.FutureTask.run(FutureTask.java:266)
        at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$201(ScheduledThreadPoolExecutor.java:180)
        at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:293)
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
        at java.lang.Thread.run(Thread.java:748)
Caused by: java.lang.IllegalStateException: A default binder has been requested, but there is more than one binder available for 'org.springframework.cloud.stream.messaging.DirectWithAttributesChannel' : kafka2,kafka1, and no default binder has been set.
        at org.springframework.cloud.stream.binder.DefaultBinderFactory.doGetBinder(DefaultBinderFactory.java:183)
        at org.springframework.cloud.stream.binder.DefaultBinderFactory.getBinder(DefaultBinderFactory.java:134)
        at org.springframework.cloud.stream.binding.BindingService.getBinder(BindingService.java:362)
        at org.springframework.cloud.stream.binding.BindingService.bindProducer(BindingService.java:257)
        at org.springframework.cloud.stream.binding.BinderAwareChannelResolver.resolveDestination(BinderAwareChannelResolver.java:125)
        at org.springframework.cloud.stream.function.FunctionConfiguration$1.lambda$afterPropertiesSet$0(FunctionConfiguration.java:188)
        at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.lang.reflect.Method.invoke(Method.java:498)
        at org.springframework.integration.handler.LambdaMessageProcessor.processMessage(LambdaMessageProcessor.java:97)
        at org.springframework.integration.router.AbstractMessageProcessingRouter.getChannelKeys(AbstractMessageProcessingRouter.java:84)
        at org.springframework.integration.router.AbstractMappingMessageRouter.determineTargetChannels(AbstractMappingMessageRouter.java:202)
        at org.springframework.integration.router.AbstractMessageRouter.handleMessageInternal(AbstractMessageRouter.java:170)
        at org.springframework.integration.handler.AbstractMessageHandler.handleMessage(AbstractMessageHandler.java:170)
        ... 26 more

Sample Project: https://github.com/ncheema/spring-cloud-stream-bug
Spring cloud stream version: 3.0.2.BUILD-SNAPSHOT

@SpringBootApplication
public class Application {
    private EmitterProcessor<Message<String>> processor = EmitterProcessor.create(false);

    @Bean
    public Supplier<Flux<Message<String>>> supplier() {
        return () -> processor;
    }

    @Bean
    public Consumer<Flux<String>> consumer() {
        return msg -> msg
                .map(m -> MessageBuilder
                        .withPayload(m)
                        .setHeader("spring.cloud.stream.sendto.destination", "topic-out")
                        .build())
                .doOnNext(processor::onNext)
                .subscribe();
    }

    public static void main(String[] args) {
        SpringApplication.run(Application.class, args);
    }
}
spring.cloud.stream:
  function.definition: supplier;consumer;
  bindings:
    supplier-out-0:
      destination: topic-out
      binder: kafka1
    consumer-in-0:
      destination: topic-in
      binder: kafka2
  binders:
    kafka1:
      type: kafka
      environment:
        spring:
          cloud:
            stream:
              kafka:
                binder:
                  brokers: localhost:9092
    kafka2:
      type: kafka
      environment:
        spring:
          cloud:
            stream:
              kafka:
                binder:
                  brokers: localhost:9092

Contributor guide

No contributing guide indexed for this repository

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.