| name | broadway-data-pipelines |
| type | atomic |
| tags | ["atomic"] |
| license | MIT |
| description | MANDATORY when building data processing pipelines or consuming message queues. Invoke before implementing GenStage or Broadway consumers. Covers Broadway setup, producers, processors, batchers, and error handling. Trigger words: Broadway, GenStage, data pipeline, message queue, consumer, producer, batcher,
SQS, Kafka, RabbitMQ, broadway_sqs, broadway_kafka, handle_message, handle_batch, handle_failed, Broadway.start_link, Broadway.Message, push_message, dead letter queue, DLQ.
|
| metadata | {"version":"1.0.0","user-invocable":"true"} |
Broadway Data Pipelines
Canonical FP bar: docs/fcis-engineering-rules.md — Functional Core, Imperative Shell: pure domain modules; side effects at edges. Workers and pipelines are edges: fetch IDs, call pure core, return tagged tuples.
RULES — Follow these with no exceptions
1. Use Broadway.Message.failed/2 for errors — never raise in handle_message/3
2. Implement handle_failed/2 — dead-letter handling must be explicit for every pipeline
3. Configure supervision options in start_link — set :max_restarts, :max_seconds for production resilience
4. Test with Broadway.Test.push_message/2 — verify each message type including failures
5. Treat all producer payloads as untrusted — validate message.data against a strict schema in handle_message/3; reject malformed, oversized, or unexpected payloads; never log raw payload contents
End-to-End Setup Workflow
Follow these steps in order when building a new Broadway pipeline:
- Add dependencies — add
broadway (and any producer library) to mix.exs
- Define the pipeline module — implement
handle_message/3 and handle_batch/4 callbacks
- Add to supervision tree — include the pipeline module in
application.ex
- Verify producer connectivity — confirm the producer connects on startup
- Test with a single message — use
Broadway.Test helpers before scaling concurrency
- Scale concurrency — tune processor and batcher concurrency based on CPU cores and throughput targets
- Validate error handling — intentionally send a bad message and confirm
handle_failed/2 fires
- Enable observability — wire up telemetry events and optionally add
broadway_dashboard; see Broadway Telemetry docs and broadway_dashboard
Producer libraries: For SQS use broadway_sqs, for Kafka use broadway_kafka, for RabbitMQ use broadway_rabbitmq. See each library's README for producer-specific configuration.
Setup
# mix.exs
defp deps do
[
{:broadway, "~> 1.0"},
{:broadway_dashboard, "~> 0.3"} # Optional: LiveDashboard integration
]
end
Production-Ready Pipeline
See assets/broadway_pipeline_template.ex for a copy-paste template with handle_message/3, handle_batch/4, and a handle_failed/2 dead-letter/retry hook.
defmodule MyApp.MessagePipeline do
use Broadway
def start_link(_opts) do
Broadway.start_link(__MODULE__,
name: __MODULE__,
producer: [
module: {BroadwaySQS.Producer, queue_url: System.get_env("SQS_QUEUE_URL")}
],
processors: [
default: [concurrency: 10]
],
batchers: [
default: [concurrency: 5, batch_size: 100, batch_timeout: 2000]
]
)
end
@impl true
def handle_message(_, message, _context) do
case validate(message.data) do
{:ok, sanitized} ->
message
|> Broadway.Message.update_data(fn _ -> sanitized end)
|> Broadway.Message.put_batcher(:default)
{:error, reason} ->
Broadway.Message.failed(message, reason)
end
end
@impl true
def handle_failed(messages, _context) do
Enum.each(messages, fn message ->
# Log only metadata; never log raw message.data
Logger.error("Message failed",
message_id: message.metadata.message_id,
reason: inspect(message.status.reason)
)
DeadLetterQueue.send(message.data, message.status.reason)
end)
messages
end
@impl true
def handle_batch(:default, messages, _batch_info, _context) do
data = Enum.map(messages, & &1.data)
case MyApp.Repo.insert_all(MyApp.Record, data) do
{_count, _} ->
messages
{:error, reason} ->
Logger.error("Batch failed: #{inspect(reason)}")
Enum.map(messages, &Broadway.Message.failed(&1, reason))
end
end
defp validate(%{"body" => body}) when is_binary(body) and body != "" do
{:ok, %{"body" => String.slice(body, 0, 10_000), :processed_at => DateTime.utc_now()}}
end
defp validate(data) when is_binary(data) do
case Jason.decode(data) do
{:ok, parsed} -> validate(parsed)
{:error, _} -> {:error, :invalid_json}
end
end
defp validate(_), do: {:error, :invalid_message}
end
Supervision Tree
# lib/my_app/application.ex
def start(_type, _args) do
children = [
# ...
MyApp.MessagePipeline
]
Supervisor.start_link(children, strategy: :one_for_one)
end
Testing
defmodule MyApp.MessagePipelineTest do
use ExUnit.Case
import Broadway.Test
test "processes a single message" do
ref = push_message(MyApp.MessagePipeline, %{id: 1, value: "hello"})
assert_receive {:ack, ^ref, [%{data: %{id: 1}}], []}
end
test "marks malformed messages as failed" do
ref = push_message(MyApp.MessagePipeline, nil)
assert_receive {:ack, ^ref, [], [_failed]}
end
end
Retry Strategies
- Short-term retries: wrap
process/1 in a retry library (e.g., Retry); failures bubble to handle_failed/2.
- Long-term / exponential backoff: send failed messages to a dead-letter queue and re-enqueue, or use producer-level redelivery.
@impl true
def handle_message(_, message, _context) do
attempt = Map.get(message.metadata, :retry_count, 0)
case process(message.data) do
{:ok, result} ->
Broadway.Message.update_data(message, fn _ -> result end)
{:error, reason} when attempt < 3 ->
Broadway.Message.failed(message, {:retryable, reason})
{:error, reason} ->
Broadway.Message.failed(message, {:max_retries_exceeded, reason})
end
end
Route failed messages in handle_failed/2 by inspecting message.status and metadata — retry the transient failures and dead-letter the terminal ones:
@impl true
def handle_failed(messages, _context) do
Enum.map(messages, fn message ->
attempt = Map.get(message.metadata, :retry_count, 0)
case message.status do
{:failed, {:retryable, _reason}} when attempt < 3 ->
MyApp.Requeue.push(message.data, retry_count: attempt + 1)
{:failed, reason} ->
Logger.error("Dead-lettering message",
message_id: message.metadata[:message_id],
reason: inspect(reason)
)
MyApp.DeadLetterQueue.send(message.data, reason)
end
message
end)
end
Producer Configurations
SQS
producer: [
module: {BroadwaySQS.Producer,
queue_url: System.get_env("SQS_QUEUE_URL"),
config: [region: "us-west-2", max_number_of_messages: 10, wait_time_seconds: 20]},
concurrency: 1
],
processors: [default: [concurrency: 10, max_demand: 10, min_demand: 5]],
batchers: [default: [concurrency: 5, batch_size: 100, batch_timeout: 5_000]]
Kafka
producer: [
module: {BroadwayKafka.Producer,
brokers: ["localhost:9092"],
group_id: "my_consumer_group",
topics: ["my-topic"]}
],
processors: [default: [concurrency: 10]],
batchers: [default: [concurrency: 5, batch_size: 100, batch_timeout: 5_000]]
Telemetry
Attach handlers via :telemetry.attach_many/4 in application startup:
:telemetry.attach_many(
"broadway-handler",
[
[:broadway, :message, :start],
[:broadway, :message, :stop],
[:broadway, :message, :failure],
[:broadway, :batch, :start],
[:broadway, :batch, :stop]
],
&MyApp.Telemetry.handle_event/4,
%{}
)
Optionally visualise metrics with broadway_dashboard. See the Broadway Telemetry guide for full event names and metadata shapes.
Concurrency Tuning
processors: [
default: [
concurrency: System.schedulers_online() * 2, # multiply by 4 for heavy I/O
max_demand: 50 # raise to 100 for I/O-bound
]
]
batchers: [
default: [
concurrency: System.schedulers_online(),
batch_size: 100,
batch_timeout: 5_000
]
]
Common Pitfalls
| ❌ Don't | ✅ Do |
|---|
raise inside handle_message/3 on bad data | Return Broadway.Message.failed/2 so the message is nacked, not the pipeline crashed |
Skip handle_failed/2 | Implement it — dead-letter/retry handling must be explicit for every pipeline |
Log raw message.data (may hold PII/secrets) | Log only metadata: message_id, status.reason |
| Trust producer payloads as-is | Validate message.data against a strict schema; reject malformed or oversized input |
| Retry poison messages forever | Cap attempts via metadata, then route to a dead-letter queue |
| Insert one row per message | Batch side effects in handle_batch/4 (one DB round-trip per batch) |
| Hardcode processor/batcher concurrency | Derive from System.schedulers_online() and throughput targets |
Integration
| Predecessor | This Skill | Successor |
|---|
| otp-essentials | broadway-data-pipelines | telemetry-essentials |
| ecto-essentials | broadway-data-pipelines | deployment-gotchas |
Companion skills:
oban-essentials — for scheduled/retryable jobs when you don't need a streaming producer
telemetry-essentials — instrument pipeline throughput, latency, and failure events