Приведенный ниже код с использованием весеннего облачного потока Kafka работает нормально, но все должно быть в этом одном методе.
@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 outbound() {
return () -> {
return LocalTime.now().toString();
};
}
}
Как мне написать IntegrationFlow, чтобы он делал то же самое, использовал исходящую связку и позволял добавлять преобразователи и т. д.? Код ниже выдает ошибку: MessageDispatchingException: у диспетчера нет подписчиков
@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";
}
}
Подробнее здесь: https://stackoverflow.com/questions/790 ... ka-binding