Skip to content

Instantly share code, notes, and snippets.

@monkey-codes
Created September 12, 2017 08:56
Show Gist options
  • Select an option

  • Save monkey-codes/c0f1762030462dcba26acaf65589ab69 to your computer and use it in GitHub Desktop.

Select an option

Save monkey-codes/c0f1762030462dcba26acaf65589ab69 to your computer and use it in GitHub Desktop.
WebFlux - Collecting user stats
public class UserStats {
UnicastProcessor eventPublisher;
Map<String, Stats> userStatsMap = new ConcurrentHashMap();
public UserStats(Flux<Event> events, UnicastProcessor eventPublisher) {
this.eventPublisher = eventPublisher;
events
.filter(type(CHAT_MESSAGE, USER_JOINED))
.subscribe(this::onChatMessage);
events
.filter(type(USER_LEFT))
.map(Event::getUser)
.map(User::getAlias)
.subscribe(userStatsMap::remove);
events
.filter(type(USER_JOINED))
.map(event -> Event.type(USER_STATS)
.withPayload()
.systemUser()
.property("stats", new HashMap<>(userStatsMap))
.build()
)
.subscribe(eventPublisher::onNext);
}
private static Predicate<Event> type(Type... types){
return event -> asList(types).contains(event.getType());
}
private void onChatMessage(Event event) {
String alias = event.getUser().getAlias();
Stats stats = userStatsMap.computeIfAbsent(alias, s -> new Stats(event.getUser()));
stats.onChatMessage(event);
}
private static class Stats {
private User user;
private long lastMessage;
private AtomicInteger messageCount = new AtomicInteger();
public Stats(User user) {
this.user = user;
}
public void onChatMessage(Event event) {
lastMessage = event.getTimestamp();
if(CHAT_MESSAGE == event.getType()) messageCount.incrementAndGet();
}
}
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment