Packages
ex_esdb
0.2.3
0.11.0
0.10.0
0.9.0
0.8.0
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.1
0.6.0
0.5.1
0.5.0
0.4.8
0.4.7
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.3
0.3.2
0.3.1
0.3.0
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
0.0.20
0.0.19
0.0.18
0.0.17
0.0.16
0.0.15
0.0.14-alpha
0.0.13-alpha
0.0.12-alpha
0.0.11-alpha
0.0.10-alpha
0.0.9-alpha
0.0.8-alpha
0.0.6-alpha
0.0.5-alpha
0.0.4-alpha
0.0.3-alpha
0.0.2-alfa
0.0.1-alfa
ExESDB is a reincarnation of rabbitmq/khepri, specialized for use as a BEAM-native event store.
Current section
Files
Jump to
Current section
Files
lib/ex_esdb/aggregator.ex
defmodule ExESDB.Aggregator do
@moduledoc """
Aggregates events from an event stream using tagged rules:
GIVEN: an Event of roughly this format:
%{
event_id: "1234567890",
event_type: "user.birthday_celebrated:v1",
stream_id: "celebrate-user-birthday-john",
version: 1,
data: %{
name: "John",
age: {:sum, 1},
venue: {:overwrite, "New York"}
},
timestamp: ~U[2022-01-01 12:00:00Z],
epoch: 1641013200,
metadata: %{
source_id: "1234567890"
}
}
"""
@doc """
Folds a list of events into a single map.
"""
def foldl(sorted_events, state \\ %{}) do
# Perform left fold (reduce)
sorted_events
|> Enum.reduce(state, fn evt, acc -> apply_evt(acc, evt) end)
end
defp get_current_num(map, key) do
case map[key] do
{:sum, value} -> value
value when is_number(value) -> value
nil -> 0
_ -> 0
end
end
defp apply_evt(state, event) do
# Process each key-value pair in the event
Map.keys(event)
|> Enum.reduce(state, fn key, acc_map ->
value = event[key]
case value do
{:sum, num} when is_number(num) ->
current = get_current_num(acc_map, key)
Map.put(acc_map, key, {:sum, current + num})
{:overwrite, new_value} ->
acc_map
|> Map.put(key, new_value)
_ ->
Map.put(acc_map, key, value)
end
end)
end
def finalize_map(tagged_map) do
Map.new(tagged_map, fn
{key, {:sum, value}} -> {key, value}
{key, {:overwrite, value}} -> {key, value}
{key, value} -> {key, value}
end)
end
end