Опросчик для отправки сообщений в привязку КафкиJAVA

Программисты JAVA общаются здесь
Anonymous
Опросчик для отправки сообщений в привязку Кафки

Сообщение Anonymous »

Приведенный ниже код с использованием весеннего облачного потока 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

Вернуться в «JAVA»