package de.juplo.kafka.chat.backend.implementation.kafka;
import de.juplo.kafka.chat.backend.ChatBackendProperties;
-import de.juplo.kafka.chat.backend.domain.ChatHomeService;
import de.juplo.kafka.chat.backend.domain.ShardingPublisherStrategy;
import de.juplo.kafka.chat.backend.implementation.haproxy.HaproxyShardingPublisherStrategy;
import de.juplo.kafka.chat.backend.implementation.kafka.messages.AbstractMessageTo;
public class KafkaServicesConfiguration
{
@Bean
- ConsumerTaskRunner consumerTaskRunner(
- ConsumerTaskExecutor infoChannelConsumerTaskExecutor,
- ConsumerTaskExecutor dataChannelConsumerTaskExecutor,
- InfoChannel infoChannel)
+ KafkaServicesThreadPoolTaskExecutorCustomizer kafkaServicesThreadPoolTaskExecutorCustomizer()
{
- return new ConsumerTaskRunner(
- infoChannelConsumerTaskExecutor,
- dataChannelConsumerTaskExecutor,
- infoChannel);
+ return new KafkaServicesThreadPoolTaskExecutorCustomizer();
}
@Bean
- ConsumerTaskExecutor infoChannelConsumerTaskExecutor(
+ ChannelTaskRunner channelTaskRunner(
+ ChannelTaskExecutor infoChannelTaskExecutor,
+ ChannelTaskExecutor dataChannelTaskExecutor)
+ {
+ return new ChannelTaskRunner(
+ infoChannelTaskExecutor,
+ dataChannelTaskExecutor);
+ }
+
+ @Bean
+ ChannelTaskExecutor infoChannelTaskExecutor(
ThreadPoolTaskExecutor taskExecutor,
InfoChannel infoChannel,
Consumer<String, AbstractMessageTo> infoChannelConsumer,
WorkAssignor infoChannelWorkAssignor)
{
- return new ConsumerTaskExecutor(
+ return new ChannelTaskExecutor(
taskExecutor,
infoChannel,
infoChannelConsumer,
}
@Bean
- ConsumerTaskExecutor dataChannelConsumerTaskExecutor(
+ ChannelTaskExecutor dataChannelTaskExecutor(
ThreadPoolTaskExecutor taskExecutor,
DataChannel dataChannel,
Consumer<String, AbstractMessageTo> dataChannelConsumer,
WorkAssignor dataChannelWorkAssignor)
{
- return new ConsumerTaskExecutor(
+ return new ChannelTaskExecutor(
taskExecutor,
dataChannel,
dataChannelConsumer,
infoChannelConsumer,
properties.getKafka().getPollingInterval(),
properties.getKafka().getNumPartitions(),
- properties.getKafka().getInstanceUri());
+ properties.getKafka().getInstanceUri(),
+ channelMediator);
channelMediator.setInfoChannel(infoChannel);
return infoChannel;
}
ChannelMediator channelMediator,
ShardingPublisherStrategy shardingPublisherStrategy)
{
- return new DataChannel(
+ DataChannel dataChannel = new DataChannel(
properties.getInstanceId(),
properties.getKafka().getDataChannelTopic(),
producer,
clock,
channelMediator,
shardingPublisherStrategy);
+ channelMediator.setDataChannel(dataChannel);
+ return dataChannel;
}
@Bean
{
return ZoneId.systemDefault();
}
+
+ @Bean
+ ChannelReactiveHealthIndicator dataChannelHealthIndicator(
+ DataChannel dataChannel)
+ {
+ return new ChannelReactiveHealthIndicator(dataChannel);
+ }
+
+ @Bean
+ ChannelReactiveHealthIndicator infoChannelHealthIndicator(InfoChannel infoChannel)
+ {
+ return new ChannelReactiveHealthIndicator(infoChannel);
+ }
}