Packages

Library to increase the throughput of producing messages (coming one at a time) to Kafka by accumulating these messages into batches

Current section

Files

Jump to
kafka_batcher lib kafka_batcher temp_storage default.ex
Raw

lib/kafka_batcher/temp_storage/default.ex

defmodule KafkaBatcher.TempStorage.Default do
@moduledoc """
Default implementation of KafkaBatcher.Behaviours.TempStorage
It just logs the messages. To have more fault tolerant implementation
you should implement your own logic to save these messages into some persistent storage.
"""
@behaviour KafkaBatcher.Behaviours.TempStorage
require Logger
@impl KafkaBatcher.Behaviours.TempStorage
def save_batch(%KafkaBatcher.TempStorage.Batch{topic: topic, partition: partition, messages: messages}) do
Logger.error("""
KafkaBatcher: Failed to send #{inspect(Enum.count(messages))} messages to the kafka topic #{topic}##{partition}
""")
Enum.each(messages, fn message -> Logger.info(inspect(message)) end)
:ok
end
@impl KafkaBatcher.Behaviours.TempStorage
def empty?(_topic) do
true
end
end