Building Storage Systems with ExCodecs

Copy Markdown View Source
# Use this install to work with the source code
# Mix.install(
#   [
#     {:ex_codecs, path: Path.join(__DIR__, "..")}, 
#     {:rustler, "~> 0.36"},
#     {:jason, "~> 1.4"}, 
#     {:kino, "~> 0.14"}, 
#     {:kino_vega_lite, "~> 0.1.13"}
#   ],
#   config:  [rustler_precompiled: [force_build: [ex_codecs: true]]]
# )

Mix.install( [
    {:ex_codecs, "~> 0.2.3"}, 
    {:jason, "~> 1.4"}, 
    {:kino, "~> 0.14"}, 
    {:kino_vega_lite, "~> 0.1.13"}
  ])

Livebook Series

#Livebook
01Introduction
02Compression Fundamentals
03Codec Comparison
04Building Storage Systems (you are here)
05Zarr-Style Workloads
06Spatial Codecs

Practical Usage Patterns

The Basic Pattern

data = "Important data to store"

{:ok, compressed} = ExCodecs.encode(:zstd, data)
{:ok, recovered} = ExCodecs.decode(:zstd, compressed)

IO.puts("Round-trip OK: #{recovered == data}")
Round-trip OK: true
:ok

Choosing a Codec at Runtime

defmodule StorageConfig do
  def codec_for_content_type("application/json"), do: :zstd
  def codec_for_content_type("application/octet-stream"), do: :lz4
  def codec_for_content_type("application/x-array"), do: :blosc2
  def codec_for_content_type(_), do: :zstd
end

content_type = "application/json"
codec = StorageConfig.codec_for_content_type(content_type)
IO.puts("Selected codec: #{inspect(codec)}")
IO.puts("Supported: #{ExCodecs.supports?(codec)}")
Selected codec: :zstd
Supported: true
:ok

Building a Compressed Storage Layer

defmodule CompressedStore do
  @moduledoc """
  A simple key-value store with transparent compression.
  """
  use Agent

  defstruct [:codec, :opts, :table]

  def start_link(opts \\ []) do
    codec = Keyword.get(opts, :codec, :zstd)
    codec_opts = Keyword.get(opts, :opts, [])
    name = Keyword.get(opts, :name, __MODULE__)

    Agent.start_link(fn ->
      %__MODULE__{codec: codec, opts: codec_opts, table: %{}}
    end, name: name)
  end

  def put(store \\ __MODULE__, key, value) do
    Agent.get_and_update(store, fn state ->
      case ExCodecs.encode(state.codec, value, state.opts) do
        {:ok, compressed} ->
          entry = %{
            compressed: compressed,
            original_size: byte_size(value),
            codec: state.codec
          }
          {{:ok, byte_size(compressed)}, %{state | table: Map.put(state.table, key, entry)}}
        {:error, err} ->
          {{:error, err}, state}
      end
    end)
  end

  def get(store \\ __MODULE__, key) do
    Agent.get(store, fn state ->
      case Map.get(state.table, key) do
        nil -> {:error, :not_found}
        entry -> ExCodecs.decode(entry.codec, entry.compressed)
      end
    end)
  end

  def info(store \\ __MODULE__) do
    Agent.get(store, fn state ->
      entries = Map.values(state.table)
      total_original = Enum.sum(Enum.map(entries, & &1.original_size))
      total_compressed = Enum.sum(Enum.map(entries, &byte_size(&1.compressed)))
      %{
        codec: state.codec,
        entry_count: length(entries),
        total_original_bytes: total_original,
        total_compressed_bytes: total_compressed,
        savings_pct: Float.round(100 * (1 - total_compressed / max(total_original, 1)), 1)
      }
    end)
  end

  def delete(store \\ __MODULE__, key) do
    Agent.get_and_update(store, fn state ->
      {Map.get(state.table, key), %{state | table: Map.delete(state.table, key)}}
    end)
  end
end
{:module, CompressedStore, <<70, 79, 82, 49, 0, 0, 36, ...>>, ...}

Using the Compressed Store

case CompressedStore.start_link(codec: :zstd, opts: [level: 3]) do
  {:ok, _pid} -> :ok
  {:error, {:already_started, _pid}} -> :ok
end

CompressedStore.put(:doc1, String.duplicate("Hello, World! ", 1000))
CompressedStore.put(:doc2, :crypto.strong_rand_bytes(4096) |> Base.encode64())
CompressedStore.put(:doc3, Jason.encode!(for i <- 1..100, do: %{id: i, value: :rand.uniform(1000)}))

{:ok, doc1} = CompressedStore.get(:doc1)
IO.puts("Retrieved doc1 (#{byte_size(doc1)} bytes)")

info = CompressedStore.info()
IO.puts("\nStorage Summary:")
IO.puts("  Codec: #{inspect(info.codec)}")
IO.puts("  Entries: #{info.entry_count}")
IO.puts("  Original: #{info.total_original_bytes} bytes")
IO.puts("  Compressed: #{info.total_compressed_bytes} bytes")
IO.puts("  Savings: #{info.savings_pct}%")
Retrieved doc1 (14000 bytes)

Storage Summary:
  Codec: :zstd
  Entries: 3
  Original: 21645 bytes
  Compressed: 4591 bytes
  Savings: 78.8%
:ok

Multi-Codec Store

defmodule MultiCodecStore do
  @moduledoc """
  A store that selects the best codec per value.
  """
  use Agent

  @codecs [:lz4, :snappy, :zstd, :bzip2]

  def start_link(opts \\ []) do
    name = Keyword.get(opts, :name, __MODULE__)
    Agent.start_link(fn -> %{} end, name: name)
  end

  def put(store \\ __MODULE__, key, value) do
    Agent.get_and_update(store, fn state ->
      results = for codec <- @codecs do
        {:ok, enc} = ExCodecs.encode(codec, value)
        {codec, enc, byte_size(enc)}
      end

      {best_codec, best_enc, best_size} =
        Enum.min_by(results, fn {_, _, size} -> size end)

      entry = %{
        codec: best_codec,
        compressed: best_enc,
        original_size: byte_size(value),
        compressed_size: best_size
      }

      {{:ok, entry}, Map.put(state, key, entry)}
    end)
  end

  def get(store \\ __MODULE__, key) do
    Agent.get(store, fn state ->
      case Map.get(state, key) do
        nil -> {:error, :not_found}
        entry -> ExCodecs.decode(entry.codec, entry.compressed)
      end
    end)
  end

  def stats(store \\ __MODULE__) do
    Agent.get(store, fn state ->
      state
      |> Map.values()
      |> Enum.group_by(& &1.codec)
      |> Enum.map(fn {codec, entries} ->
        {codec, length(entries),
         Enum.sum(Enum.map(entries, & &1.original_size)),
         Enum.sum(Enum.map(entries, & &1.compressed_size))}
      end)
      |> Enum.sort_by(&elem(&1, 1))
    end)
  end
end
{:module, MultiCodecStore, <<70, 79, 82, 49, 0, 0, 27, ...>>, ...}

Using the Multi-Codec Store

case MultiCodecStore.start_link() do
  {:ok, _pid} -> :ok
  {:error, {:already_started, _pid}} -> :ok
end

datasets = %{
  repetitive: String.duplicate("abcabcabc", 5000),
  semi_random: (for _ <- 1..10000, into: <<>>, do: <<:rand.uniform(255)>>),
  floats: (for i <- 1..4096, into: <<>>, do: <<i * 0.1::float-size(64)-little>>)
}

for {name, data} <- datasets do
  {:ok, entry} = MultiCodecStore.put(name, data)
  IO.puts("#{name}: best codec = #{inspect(entry.codec)}, " <>
    "#{entry.original_size} -> #{entry.compressed_size} bytes " <>
    "(#{Float.round(100 * entry.compressed_size / entry.original_size, 1)}%)")
end

IO.puts("\nCodec distribution:")
for {codec, count, orig, comp} <- MultiCodecStore.stats() do
  IO.puts("  #{inspect(codec)}: #{count} entries, #{orig} -> #{comp} bytes")
end
floats: best codec = :zstd, 32768 -> 4801 bytes (14.7%)
repetitive: best codec = :zstd, 45000 -> 21 bytes (0.0%)
semi_random: best codec = :snappy, 10000 -> 10005 bytes (100.0%)

Codec distribution:
  :snappy: 1 entries, 10000 -> 10005 bytes
  :zstd: 2 entries, 77768 -> 4822 bytes
[:ok, :ok]

Chunked Compression for Large Files

defmodule ChunkedCompressor do
  @moduledoc """
  Compresses large data in fixed-size chunks to limit peak memory usage.

  The container stores the chunk count, the original byte size, and an offset
  table so individual chunks can be decompressed on demand (random access)
  without re-decoding the whole blob.
  """
  @default_chunk_size 64 * 1024

  def compress(data, codec \\ :zstd, opts \\ [], chunk_size \\ @default_chunk_size) do
    chunks = chunk_data(data, chunk_size)

    compressed_chunks =
      Enum.map(chunks, fn chunk ->
        {:ok, enc} = ExCodecs.encode(codec, chunk, opts)
        enc
      end)

    offsets =
      compressed_chunks
      |> Enum.scan(0, fn chunk, acc -> acc + byte_size(chunk) end)
      |> Enum.drop(-1)
      |> List.insert_at(0, 0)

    header = <<
      length(compressed_chunks)::unsigned-big-32,
      byte_size(data)::unsigned-big-64
    >>

    offset_data = for offset <- offsets, into: <<>>, do: <<offset::unsigned-big-32>>

    {:ok, header <> offset_data <> IO.iodata_to_binary(compressed_chunks)}
  end

  def decompress(blob, codec \\ :zstd) do
    <<num_chunks::unsigned-big-32, original_size::unsigned-big-64, rest::binary>> = blob

    offsets_size = num_chunks * 4
    <<offsets::binary-size(^offsets_size), chunks_binary::binary>> = rest

    offsets_list = for <<offset::unsigned-big-32 <- offsets>>, do: offset
    offsets_list = offsets_list ++ [byte_size(chunks_binary)]

    chunks =
      Enum.chunk_every(offsets_list, 2, 1, :discard)
      |> Enum.map(fn [start_off, end_off] ->
        binary_part(chunks_binary, start_off, end_off - start_off)
      end)

    decompressed =
      Enum.map(chunks, fn chunk ->
        {:ok, dec} = ExCodecs.decode(codec, chunk)
        dec
      end)

    result = IO.iodata_to_binary(decompressed)

    if byte_size(result) == original_size do
      {:ok, result}
    else
      {:error, :size_mismatch}
    end
  end

  @doc """
  Decompresses a single chunk by 0-based index, using the offset table.
  """
  def read_chunk(blob, index, codec \\ :zstd) do
    <<num_chunks::unsigned-big-32, _original_size::unsigned-big-64, rest::binary>> = blob

    if index < 0 or index >= num_chunks do
      {:error, :out_of_range}
    else
      offsets_size = num_chunks * 4
      <<offsets::binary-size(^offsets_size), chunks_binary::binary>> = rest

      offsets_list = for <<offset::unsigned-big-32 <- offsets>>, do: offset
      offsets_list = offsets_list ++ [byte_size(chunks_binary)]

      [start_off, end_off] = Enum.at(Enum.chunk_every(offsets_list, 2, 1, :discard), index)
      chunk = binary_part(chunks_binary, start_off, end_off - start_off)

      ExCodecs.decode(codec, chunk)
    end
  end

  defp chunk_data(data, chunk_size) when is_integer(chunk_size) and chunk_size > 0 do
    case data do
      <<chunk::binary-size(^chunk_size), rest::binary>> -> [chunk | chunk_data(rest, chunk_size)]
      <<>> -> []
      remainder when is_binary(remainder) -> [remainder]
    end
  end
end
{:module, ChunkedCompressor, <<70, 79, 82, 49, 0, 0, 32, ...>>, ...}
large_data = String.duplicate("This is a chunk of data that repeats. ", 10000)

{:ok, compressed} = ChunkedCompressor.compress(large_data, :zstd)
{:ok, recovered} = ChunkedCompressor.decompress(compressed, :zstd)

<<num_chunks::unsigned-big-32, _::unsigned-big-64, _::binary>> = compressed
{:ok, first_chunk} = ChunkedCompressor.read_chunk(compressed, 0, :zstd)

IO.puts("Original:   #{byte_size(large_data)} bytes")
IO.puts("Compressed: #{byte_size(compressed)} bytes")
IO.puts("Ratio:      #{Float.round(100 * byte_size(compressed) / byte_size(large_data), 1)}%")
IO.puts("Chunks:     #{num_chunks}")
IO.puts("Round-trip: #{recovered == large_data}")
IO.puts("Random access chunk 0: #{byte_size(first_chunk)} bytes (matches: #{first_chunk == binary_part(large_data, 0, 64 * 1024)})")
Original:   380000 bytes
Compressed: 384 bytes
Ratio:      0.1%
Chunks:     6
Round-trip: true
Random access chunk 0: 65536 bytes (matches: true)
:ok

Error Handling Best Practices

case ExCodecs.encode(:zstd, "some data") do
  {:ok, compressed} ->
    {:ok, compressed}

  {:error, %ExCodecs.Error{reason: :unsupported_codec}} ->
    {:error, :codec_missing}

  {:error, %ExCodecs.Error{reason: :codec_unavailable}} ->
    {:error, :codec_unavailable}

  {:error, %ExCodecs.Error{reason: :nif_not_loaded}} ->
    {:error, :nif_not_loaded}

  {:error, %ExCodecs.Error{reason: :invalid_data} = err} ->
    IO.puts("Invalid data: #{err.message}")
    {:error, :bad_input}

  {:error, %ExCodecs.Error{reason: :compression_failed} = err} ->
    IO.puts("Compression failed: #{err.message}")
    {:error, :compress_failed}

  {:error, %ExCodecs.Error{} = err} ->
    IO.puts("Structured error #{err.reason}: #{err.message}")
    {:error, err.reason}
end
{:ok, <<40, 181, 47, 253, 32, 9, 73, 0, 0, 115, 111, 109, 101, 32, 100, 97, 116, 97>>}
# Exercise the error branches with real failures.
results =
  [
    {:ok, _} = ExCodecs.encode(:zstd, "ok"),
    {:error, %ExCodecs.Error{reason: :invalid_options}} =
      ExCodecs.encode(:zstd, "x", level: 99),
    {:error, %ExCodecs.Error{reason: :decompression_failed}} =
      ExCodecs.decode(:zstd, "not a zstd frame"),
    {:error, %ExCodecs.Error{reason: :unsupported_codec}} =
      ExCodecs.encode(:nonexistent, "x")
  ]

for {tag, {status, _} = result} <- Enum.zip([:ok, :bad_level, :bad_frame, :bad_codec], results) do
  IO.puts("#{String.pad_trailing("#{tag}", 12)} -> #{inspect(result)}")
end
warning: variable "status" is unused (if the variable is not meant to be used, prefix it with an underscore)
 work/ex_codecs/livebooks/04_building_storage_systems.livemd#cell:pruxa6uyvvgiepiz:13

ok           -> {:ok, <<40, 181, 47, 253, 32, 2, 17, 0, 0, 111, 107>>}
bad_level    -> {:error, %ExCodecs.Error{reason: :invalid_options, message: "Level must be an integer between 1 and 22", codec: nil, details: nil}}
bad_frame    -> {:error, %ExCodecs.Error{reason: :decompression_failed, message: "Decompression failed", codec: :zstd, details: nil}}
bad_codec    -> {:error, %ExCodecs.Error{reason: :unsupported_codec, message: "The specified codec is not supported", codec: :nonexistent, details: nil}}
[:ok, :ok, :ok, :ok]
# Bound decompression for untrusted payloads (default is 256 MiB).
# Classic bomb shape: 64 KiB of zeros → tiny compressed input.
bomb_raw = :binary.copy(<<0>>, 65_536)
{:ok, bomb} = ExCodecs.encode(:zstd, bomb_raw)

{:error, %ExCodecs.Error{reason: :output_limit_exceeded}} =
  ExCodecs.decode(:zstd, bomb, max_output_size: 1_024)

{:ok, ^bomb_raw} = ExCodecs.decode(:zstd, bomb, max_output_size: 65_536)

IO.puts("bomb compressed size: #{byte_size(bomb)} bytes (#{Float.round(65_536 / byte_size(bomb), 1)}x)")
bomb compressed size: 11 bytes (5957.8x)
:ok

See Compression Fundamentals for the same pattern across all compression codecs.

defmodule SafeCompress do
  def call(codec, data, opts \\ []) do
    with {:check_codec, true} <- {:check_codec, ExCodecs.supports?(codec)},
         {:check_binary, true} <- {:check_binary, is_binary(data)},
         {:ok, compressed} <- ExCodecs.encode(codec, data, opts),
         {:ok, ^data} <- ExCodecs.decode(codec, compressed) do
      {:ok, compressed}
    else
      {:check_codec, false} -> {:error, :unsupported_codec}
      {:check_binary, false} -> {:error, :not_binary}
      {:error, err} -> {:error, err}
    end
  end
end

{:ok, result} = SafeCompress.call(:zstd, "verify round-trip")
IO.puts("Safe compress OK: #{is_binary(result)}")
Safe compress OK: true
:ok
defmodule CompressWithFallback do
  def call(data, codecs \\ [:zstd, :lz4, :snappy]) do
    codecs
    |> Enum.reduce_while({:error, :no_codec_available}, fn codec, _acc ->
      case ExCodecs.encode(codec, data) do
        {:ok, compressed} -> {:halt, {:ok, {codec, compressed}}}
        {:error, err} -> {:cont, {:error, {codec, err.reason}}}
      end
    end)
  end
end

{:ok, {used_codec, compressed}} = CompressWithFallback.call("fallback test data")
IO.puts("Used codec: #{inspect(used_codec)}")
IO.puts("Compressed: #{byte_size(compressed)} bytes")
Used codec: :zstd
Compressed: 27 bytes
:ok

ETS Integration

defmodule CompressedETS do
  @moduledoc """
  ETS-backed storage with transparent compression.
  """
  def new(name \\ :compressed_store) do
    if :ets.whereis(name) != :undefined, do: :ets.delete(name)
    :ets.new(name, [:set, :public, :named_table])
    name
  end

  def insert(table, key, value, codec \\ :zstd, opts \\ []) do
    {:ok, compressed} = ExCodecs.encode(codec, value, opts)
    :ets.insert(table, {key, compressed, codec, byte_size(value)})
    {:ok, byte_size(compressed)}
  end

  def lookup(table, key) do
    case :ets.lookup(table, key) do
      [{^key, compressed, codec, _orig_size}] ->
        ExCodecs.decode(codec, compressed)
      [] ->
        {:error, :not_found}
    end
  end

  def stats(table) do
    entries = :ets.tab2list(table)
    total_orig = Enum.sum(Enum.map(entries, fn {_, _, _, orig} -> orig end))
    total_comp = Enum.sum(Enum.map(entries, fn {_, comp, _, _} -> byte_size(comp) end))
    %{
      entries: length(entries),
      original_bytes: total_orig,
      compressed_bytes: total_comp,
      savings_pct: Float.round(100 * (1 - total_comp / max(total_orig, 1)), 1)
    }
  end
end
{:module, CompressedETS, <<70, 79, 82, 49, 0, 0, 20, ...>>, ...}
table = CompressedETS.new()

for i <- 1..100 do
  data = String.duplicate("item-#{rem(i, 10)} ", 500)
  CompressedETS.insert(table, "key_#{i}", data)
end

stats = CompressedETS.stats(table)
IO.puts("ETS Compressed Store:")
IO.puts("  Entries: #{stats.entries}")
IO.puts("  Original: #{stats.original_bytes} bytes")
IO.puts("  Compressed: #{stats.compressed_bytes} bytes")
IO.puts("  Savings: #{stats.savings_pct}%")

{:ok, val} = CompressedETS.lookup(table, "key_1")
IO.puts("\nRetrieved #{byte_size(val)} bytes successfully")

:ets.delete(table)
ETS Compressed Store:
  Entries: 100
  Original: 350000 bytes
  Compressed: 2400 bytes
  Savings: 99.3%

Retrieved 3500 bytes successfully
true

File Storage Integration

defmodule CompressedFileIO do
  @moduledoc """
  File I/O with transparent compression.
  """
  @magic "EXCD"

  def write(path, data, codec \\ :zstd, opts \\ []) do
    {:ok, compressed} = ExCodecs.encode(codec, data, opts)
    header = <<
      @magic::binary,
      codec_atom_to_bin(codec)::binary,
      byte_size(data)::unsigned-big-64,
      byte_size(compressed)::unsigned-big-64
    >>
    File.write(path, header <> compressed)
  end

  def read(path) do
    case File.read(path) do
      {:ok, <<@magic::binary, codec_bin::binary-size(8), _orig_size::unsigned-big-64, _comp_size::unsigned-big-64, compressed::binary>>} ->
        codec = bin_to_codec_atom(codec_bin)
        ExCodecs.decode(codec, compressed)
      {:ok, _} ->
        {:error, :invalid_format}
      {:error, reason} ->
        {:error, reason}
    end
  end

  defp codec_atom_to_bin(atom) do
    atom |> Atom.to_string() |> String.pad_trailing(8)
  end

  defp bin_to_codec_atom(bin) do
    case bin |> String.trim_trailing() do
      "zstd" -> :zstd
      "lz4" -> :lz4
      "snappy" -> :snappy
      "bzip2" -> :bzip2
      "blosc2" -> :blosc2
      other -> {:error, {:unknown_codec, other}}
    end
  end
end
{:module, CompressedFileIO, <<70, 79, 82, 49, 0, 0, 16, ...>>, ...}
test_data = String.duplicate("Compressed file storage example. ", 2000)
path = Path.join(System.tmp_dir!(), "ex_codecs_demo.dat")

:ok = CompressedFileIO.write(path, test_data, :zstd, level: 9)
{:ok, recovered} = CompressedFileIO.read(path)

IO.puts("Wrote and read #{byte_size(test_data)} bytes")
IO.puts("Round-trip OK: #{recovered == test_data}")
IO.puts("File size: #{File.stat!(path).size} bytes")

File.rm(path)
Wrote and read 66000 bytes
Round-trip OK: true
File size: 83 bytes
:ok

Key Takeaways

  1. Wrap encode/decode in your own modules to centralize codec choice and error handling
  2. Always handle errors — codecs can fail on invalid input, memory pressure, etc.
  3. Store the codec name alongside compressed data so you know how to decompress
  4. Use ETS for compressed in-memory caches — significant memory savings
  5. Chunk large files — enables random access, parallel processing, and bounded memory
  6. Try multiple codecs for mixed workloads — different data compresses best with different algorithms

Previous: Codec Comparison · Next: Zarr-Style Workloads