Partition routing helpers for multi-partition streams.
Routing uses :erlang.phash2/2 (stable on the BEAM; not portable to other
runtimes). Prefer an explicit :partition when you need a fixed assignment.
append(payload, key: "EUR")
|
v
:erlang.phash2(key, N) --> partition in 0..N-1
|
v
per-partition sequence ++ immutable event
Summary
Functions
String key for maps stored in Mongo ("0", "1", …).
Reads an integer from a string- or integer-keyed map (Mongo/ETS frontiers).
Normalizes :partition / :partitions / :all into a sorted unique list.
Resolves the target partition for an append.
Returns {:ok, p} when p is in 0..partition_count-1.
Functions
@spec key(non_neg_integer()) :: String.t()
String key for maps stored in Mongo ("0", "1", …).
Examples
iex> Rheo.Partition.key(0)
"0"
@spec map_get(map(), non_neg_integer(), non_neg_integer()) :: non_neg_integer()
Reads an integer from a string- or integer-keyed map (Mongo/ETS frontiers).
Examples
iex> Rheo.Partition.map_get(%{"0" => 5}, 0, 0)
5
iex> Rheo.Partition.map_get(%{0 => 3}, 0, 0)
3
iex> Rheo.Partition.map_get(%{}, 1, 0)
0
@spec normalize_assignment(term(), pos_integer()) :: {:ok, [non_neg_integer()]} | {:error, :invalid_partition}
Normalizes :partition / :partitions / :all into a sorted unique list.
Used by Group/Producer assignment opts.
Examples
iex> Rheo.Partition.normalize_assignment(:all, 3)
{:ok, [0, 1, 2]}
iex> Rheo.Partition.normalize_assignment([2, 0, 2], 4)
{:ok, [0, 2]}
iex> Rheo.Partition.normalize_assignment(1, 4)
{:ok, [1]}
iex> Rheo.Partition.normalize_assignment([9], 4)
{:error, :invalid_partition}
@spec resolve(map(), keyword(), pos_integer()) :: {:ok, non_neg_integer()} | {:error, :invalid_partition}
Resolves the target partition for an append.
Precedence: explicit :partition in opts, else :key in opts or payload
(:key / "key"), else 0.
Arguments
payload— event body map (may contain:key/"key")opts— keyword list; recognized keys:partition,:keypartition_count— stream partition count (>= 1)
Returns
{:ok, partition}when in0..partition_count-1{:error, :invalid_partition}when:partitionis out of range
Examples
iex> Rheo.Partition.resolve(%{}, [], 4)
{:ok, 0}
iex> Rheo.Partition.resolve(%{}, [partition: 2], 4)
{:ok, 2}
iex> Rheo.Partition.resolve(%{}, [partition: 9], 4)
{:error, :invalid_partition}
iex> {:ok, p} = Rheo.Partition.resolve(%{"key" => "EUR"}, [], 4)
iex> p in 0..3
true
@spec validate(term(), pos_integer()) :: {:ok, non_neg_integer()} | {:error, :invalid_partition}
Returns {:ok, p} when p is in 0..partition_count-1.
Examples
iex> Rheo.Partition.validate(0, 4)
{:ok, 0}
iex> Rheo.Partition.validate(4, 4)
{:error, :invalid_partition}