Packages
kafka_ex_tc
0.12.1-21
0.13.0
0.12.1
0.12.1-51-1
0.12.1-50-3
0.12.1-50-2
0.12.1-50-1
0.12.1-49-2
0.12.1-49-1
0.12.1-49-0
0.12.1-48-2
0.12.1-48-1
0.12.1-47-1
0.12.1-46-2
0.12.1-45-2
0.12.1-45-1
0.12.1-44-1
0.12.1-43-1
0.12.1-42-8
0.12.1-42-6
0.12.1-42-5
0.12.1-42-4
0.12.1-42-1
0.12.1-41-2
0.12.1-41-1
0.12.1-40-9
0.12.1-40-8
0.12.1-40-7
0.12.1-40-6
0.12.1-40-5
0.12.1-40-4
0.12.1-40-3
0.12.1-40-2
0.12.1-40-10
0.12.1-40-1
0.12.1-39-2
0.12.1-39-1
0.12.1-36-3
0.12.1-36-2
0.12.1-36-1
0.12.1-34-b
0.12.1-34-a
0.12.1-33-1
0.12.1-32-a1
0.12.1-29-s-3
0.12.1-29-s-2
0.12.1-29-s-1
0.12.1-28-rumble-1
0.12.1-27-2
0.12.1-27-1
0.12.1-26-rumble-2
0.12.1-26-rumble-1
0.12.1-26-dev-1
0.12.1-26-dev.1
0.12.1-26-3
0.12.1-26-2
0.12.1-26-1
0.12.1-25-s-8
0.12.1-25-s-7
0.12.1-25-s-6
0.12.1-25-s-5
0.12.1-25-s-4
0.12.1-25-s-3
0.12.1-25-s-2
0.12.1-25-s-1
0.12.1-25-dev-22
0.12.1-25-dev
0.12.1-24-2
0.12.1-24-1
0.12.1-22-poc-b-7
0.12.1-20-poc-b-6
0.12.1-20-poc-b-5
0.12.1-20-poc-b-4
0.12.1-20-poc-b-3
0.12.1-19-poc-b-3
0.12.1-19-poc-b-2
0.12.1-19-poc-b-1
0.12.1-19-poc-7
0.12.1-19-poc-6
0.12.1-19-poc-4
0.12.1-19-poc-3
0.12.1-19-poc-2
0.12.1-19-poc-1
0.12.1-19-poc
0.12.1-11-debug-1
0.12.1-11-debug
0.12.1-39
0.12.1-38.a
0.12.1-38
0.12.1-37.c
0.12.1-37.b
0.12.1-37.a
0.12.1-37.1
0.12.1-37
0.12.1-35.b
0.12.1-35.a
0.12.1-34
0.12.1-33
0.12.1-32
0.12.1-31
0.12.1-30
0.12.1-29
0.12.1-28
0.12.1-27
0.12.1-26
0.12.1-25
0.12.1-24
0.12.1-23
0.12.1-22.1
0.12.1-22
0.12.1-21
0.12.1-20
0.12.1-19
0.12.1-18
0.12.1-17
0.12.1-16
0.12.1-15
0.12.1-14
0.12.1-13
0.12.1-12
0.12.1-11
0.12.1-10
0.12.1-9
0.12.1-8
0.12.1-7
0.12.1-6
0.12.1-5
0.12.1-4
0.12.1-3
0.12.1-2
0.12.1-1
Kafka client for Elixir/Erlang.
Current section
Files
Jump to
Current section
Files
lib/kafka_ex/compression.ex
defmodule KafkaEx.Compression do
@moduledoc """
Handles compression/decompression of messages.
See https://cwiki.apache.org/confluence/display/KAFKA/A+Guide+To+The+Kafka+Protocol#AGuideToTheKafkaProtocol-Compression
To add new compression types:
1. Add the appropriate dependency to mix.exs (don't forget to add it
to the application list).
2. Add the appropriate attribute value and compression_type atom.
3. Add a decompress function clause.
4. Add a compress function clause.
"""
@gzip_attribute 1
@snappy_attribute 2
@type attribute_t :: integer
@type compression_type_t :: :snappy | :gzip
@doc """
This function should pattern match on the attribute value and return
the decompressed data.
"""
@spec decompress(attribute_t, binary) :: binary
def decompress(@gzip_attribute, data) do
:zlib.gunzip(data)
end
def decompress(@snappy_attribute, data) do
<<_snappy_header::64, _snappy_version_info::64, rest::binary>> = data
snappy_decompress_chunk(rest, <<>>)
end
@doc """
This function should pattern match on the compression_type atom and
return the compressed data as well as the corresponding attribute
value.
"""
@spec compress(compression_type_t, binary) :: {binary, attribute_t}
def compress(:snappy, data) do
{:ok, compressed_data} = :snappy.compress(data)
{compressed_data, @snappy_attribute}
end
def compress(:gzip, data) do
compressed_data = :zlib.gzip(data)
{compressed_data, @gzip_attribute}
end
def snappy_decompress_chunk(<<>>, so_far) do
so_far
end
def snappy_decompress_chunk(
<<valsize::32-unsigned, value::size(valsize)-binary, rest::binary>>,
so_far
) do
{:ok, decompressed_value} = :snappy.decompress(value)
snappy_decompress_chunk(rest, so_far <> decompressed_value)
end
end