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;
import de.juplo.kafka.chat.backend.implementation.kafka.messages.data.EventChatMessageReceivedTo;
import de.juplo.kafka.chat.backend.implementation.kafka.messages.info.EventChatRoomCreated;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
+import java.net.InetSocketAddress;
import java.time.Clock;
import java.time.ZoneId;
import java.util.HashMap;
public class KafkaServicesConfiguration
{
@Bean
- ConsumerTaskRunner consumerTaskRunner(
- ConsumerTaskExecutor infoChannelConsumerTaskExecutor,
- ConsumerTaskExecutor dataChannelConsumerTaskExecutor)
+ KafkaServicesThreadPoolTaskExecutorCustomizer kafkaServicesThreadPoolTaskExecutorCustomizer()
{
- return new ConsumerTaskRunner(
- infoChannelConsumerTaskExecutor,
- dataChannelConsumerTaskExecutor);
+ return new KafkaServicesThreadPoolTaskExecutorCustomizer();
}
- @Bean
- ConsumerTaskExecutor infoChannelConsumerTaskExecutor(
+ @Bean(initMethod = "executeChannelTask", destroyMethod = "join")
+ ChannelTaskExecutor infoChannelTaskExecutor(
ThreadPoolTaskExecutor taskExecutor,
InfoChannel infoChannel,
Consumer<String, AbstractMessageTo> infoChannelConsumer,
WorkAssignor infoChannelWorkAssignor)
{
- return new ConsumerTaskExecutor(
+ return new ChannelTaskExecutor(
taskExecutor,
infoChannel,
infoChannelConsumer,
};
}
- @Bean
- ConsumerTaskExecutor dataChannelConsumerTaskExecutor(
+ @Bean(initMethod = "executeChannelTask", destroyMethod = "join")
+ ChannelTaskExecutor dataChannelTaskExecutor(
ThreadPoolTaskExecutor taskExecutor,
DataChannel dataChannel,
Consumer<String, AbstractMessageTo> dataChannelConsumer,
WorkAssignor dataChannelWorkAssignor)
{
- return new ConsumerTaskExecutor(
+ return new ChannelTaskExecutor(
taskExecutor,
dataChannel,
dataChannelConsumer,
}
@Bean
- ChatHomeService kafkaChatHome(
+ KafkaChatHomeService kafkaChatHome(
ChatBackendProperties properties,
InfoChannel infoChannel,
DataChannel dataChannel)
InfoChannel infoChannel(
ChatBackendProperties properties,
Producer<String, AbstractMessageTo> producer,
- Consumer<String, AbstractMessageTo> infoChannelConsumer)
+ Consumer<String, AbstractMessageTo> infoChannelConsumer,
+ ChannelMediator channelMediator)
{
- return new InfoChannel(
+ InfoChannel infoChannel = new InfoChannel(
properties.getKafka().getInfoChannelTopic(),
producer,
- infoChannelConsumer);
+ infoChannelConsumer,
+ properties.getKafka().getPollingInterval(),
+ properties.getKafka().getNumPartitions(),
+ properties.getKafka().getInstanceUri(),
+ channelMediator);
+ channelMediator.setInfoChannel(infoChannel);
+ return infoChannel;
}
@Bean
Consumer<String, AbstractMessageTo> dataChannelConsumer,
ZoneId zoneId,
Clock clock,
- InfoChannel infoChannel)
+ ChannelMediator channelMediator,
+ ShardingPublisherStrategy shardingPublisherStrategy)
{
- return new DataChannel(
+ DataChannel dataChannel = new DataChannel(
+ properties.getInstanceId(),
properties.getKafka().getDataChannelTopic(),
producer,
dataChannelConsumer,
zoneId,
properties.getKafka().getNumPartitions(),
+ properties.getKafka().getPollingInterval(),
properties.getChatroomBufferSize(),
clock,
- infoChannel);
+ channelMediator,
+ shardingPublisherStrategy);
+ channelMediator.setDataChannel(dataChannel);
+ return dataChannel;
+ }
+
+ @Bean
+ ChannelMediator channelMediator()
+ {
+ return new ChannelMediator();
}
@Bean
return properties;
}
+ @Bean
+ ShardingPublisherStrategy shardingPublisherStrategy(
+ ChatBackendProperties properties)
+ {
+ String[] parts = properties.getKafka().getHaproxyRuntimeApi().split(":");
+ InetSocketAddress haproxyAddress = new InetSocketAddress(parts[0], Integer.valueOf(parts[1]));
+ return new HaproxyShardingPublisherStrategy(
+ haproxyAddress,
+ properties.getKafka().getHaproxyMap(),
+ properties.getInstanceId());
+ }
+
@Bean
ZoneId zoneId()
{
return ZoneId.systemDefault();
}
+
+ @Bean
+ ChannelReactiveHealthIndicator dataChannelHealthIndicator(
+ DataChannel dataChannel)
+ {
+ return new ChannelReactiveHealthIndicator(dataChannel);
+ }
+
+ @Bean
+ ChannelReactiveHealthIndicator infoChannelHealthIndicator(InfoChannel infoChannel)
+ {
+ return new ChannelReactiveHealthIndicator(infoChannel);
+ }
}