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.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.JsonDeserializer;
import org.springframework.kafka.support.serializer.JsonSerializer;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
+import reactor.core.publisher.Mono;
import java.time.Clock;
import java.time.ZoneId;
@Bean
ConsumerTaskRunner consumerTaskRunner(
ConsumerTaskExecutor infoChannelConsumerTaskExecutor,
- ConsumerTaskExecutor dataChannelConsumerTaskExecutor)
+ ConsumerTaskExecutor dataChannelConsumerTaskExecutor,
+ InfoChannel infoChannel)
{
return new ConsumerTaskRunner(
infoChannelConsumerTaskExecutor,
- dataChannelConsumerTaskExecutor);
+ dataChannelConsumerTaskExecutor,
+ infoChannel);
}
@Bean
ThreadPoolTaskExecutor taskExecutor,
InfoChannel infoChannel,
Consumer<String, AbstractMessageTo> infoChannelConsumer,
- ConsumerTaskExecutor.WorkAssignor infoChannelWorkAssignor)
+ WorkAssignor infoChannelWorkAssignor)
{
return new ConsumerTaskExecutor(
taskExecutor,
}
@Bean
- ConsumerTaskExecutor.WorkAssignor infoChannelWorkAssignor(
- ChatBackendProperties properties)
+ WorkAssignor infoChannelWorkAssignor(ChatBackendProperties properties)
{
return consumer ->
{
ThreadPoolTaskExecutor taskExecutor,
DataChannel dataChannel,
Consumer<String, AbstractMessageTo> dataChannelConsumer,
- ConsumerTaskExecutor.WorkAssignor dataChannelWorkAssignor)
+ WorkAssignor dataChannelWorkAssignor)
{
return new ConsumerTaskExecutor(
taskExecutor,
}
@Bean
- ConsumerTaskExecutor.WorkAssignor dataChannelWorkAssignor(
+ WorkAssignor dataChannelWorkAssignor(
ChatBackendProperties properties,
DataChannel dataChannel)
{
return new InfoChannel(
properties.getKafka().getInfoChannelTopic(),
producer,
- infoChannelConsumer);
+ infoChannelConsumer,
+ properties.getKafka().getInstanceUri());
}
@Bean
Producer<String, AbstractMessageTo> producer,
Consumer<String, AbstractMessageTo> dataChannelConsumer,
ZoneId zoneId,
- Clock clock)
+ Clock clock,
+ InfoChannel infoChannel,
+ ShardingPublisherStrategy shardingPublisherStrategy)
{
return new DataChannel(
+ properties.getInstanceId(),
properties.getKafka().getDataChannelTopic(),
producer,
dataChannelConsumer,
zoneId,
properties.getKafka().getNumPartitions(),
properties.getChatroomBufferSize(),
- clock);
+ clock,
+ infoChannel,
+ shardingPublisherStrategy);
}
@Bean
return properties;
}
+ @Bean
+ ShardingPublisherStrategy shardingPublisherStrategy()
+ {
+ return new ShardingPublisherStrategy() {
+ @Override
+ public Mono<String> publishOwnership(int shard)
+ {
+ return Mono.just(Integer.toString(shard));
+ }
+ };
+ }
+
@Bean
ZoneId zoneId()
{