Rheo.Partition (rheo v1.0.0)

Copy Markdown View Source

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

key(partition)

@spec key(non_neg_integer()) :: String.t()

String key for maps stored in Mongo ("0", "1", …).

Examples

iex> Rheo.Partition.key(0)
"0"

map_get(map, partition, default \\ 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

normalize_assignment(partition, partition_count)

@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}

resolve(payload, opts, partition_count)

@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, :key
  • partition_count — stream partition count (>= 1)

Returns

  • {:ok, partition} when in 0..partition_count-1
  • {:error, :invalid_partition} when :partition is 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

validate(partition, partition_count)

@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}