Packages

Amazon Kinesis Data Firehose configurable queue supporting arbitrary adapters

Current section

Files

Jump to
firefighter lib adapters ex_aws.ex
Raw

lib/adapters/ex_aws.ex

if Code.ensure_loaded?(ExAws.Firehose) do
defmodule Firefighter.Adapters.ExAws do
require Logger
@behaviour Firefighter.Adapter
def pump(stream_name, records, _opts) do
record_batch =
records
|> Enum.map(fn record ->
record <> delimiter()
end)
result =
ExAws.Firehose.put_record_batch(stream_name, record_batch)
|> ExAws.request()
case result do
{:ok, response} -> {:ok, response}
{:error, error} -> {:error, error}
error -> {:error, error}
end
end
defp delimiter, do: Application.get_env(:firefighter, :delimiter, "")
end
end