giddie
Elixir-postgresql-message-queue - Pure PostgreSQL Message Queue for Elixir
Phoenix Pub/Sub is great, but I often encounter usecases where I want topic-based message passing, but durable. In other words, if the message is sent, I want to guarantee that it will eventually hit all of its configured listeners.
Oban is very popular, but the “job” abstraction feels too heavy for me. RabbitMQ+Broadway is a great combo, but it adds complexity to the deployment.
So for one project I decided to implement a reasonably flexible message queue system using only PostgreSQL. I think I’ll want to use this again, so I’ve distilled out the relevant code into a reference repo:
Usage looks a bit like this:
iex> Messaging.MessageQueueProcessor.start_link(queue: "my_queue")
iex> [%Message{type: "Example.Event", schema_version: 1, payload: %{"one" => 1}}]
...> |> Messaging.broadcast_messages!(to_queue: "my_queue")
config :postgresql_message_queue, PostgresqlMessageQueue.Messaging,
broadcast_listeners: [
{MyContext.MyMessageHandler, ["MyContext.Commands.*", "AnotherContext.Events.*"]},
{MyLogger.EventLogger, ["*.Events.*"]}
]
@impl Messaging.MessageHandler
def handle_message(%Messaging.Message{
type: "ExampleUsage.Events.Greeting",
payload: %{"greeting" => greeting}
}) do
Logger.info("ExampleUsage: received greeting: #{greeting}")
end
I’d welcome any feedback, and I hope it’s useful to someone else, or at least interesting. I’m still unsure if this would work well as a library, but I’m thinking about it.
Most Liked
benwilson512
Hey @giddie! What guarantees does this project provide regarding message order and visibility? That is to say, one of the common “gotchas” with Postgres is assuming that things like sequences can be relied upon to act like monotonic cursors. For example:
Process A | Process B
BEGIN | BEGIN
insert next_val('my_seq') | insert next_val('my_seq')
1 | 2
# do something that takes time | COMMIT
COMMIT
In this scenario, the value 2 will be visible outside the transaction before 1, and so if a consumer treats “oh I have seen message with id 2, therefore I am caught up to 2” they will miss messages.
How does this project avoid that issue?
r8code
looks good
there also GitHub - gmtprime/yggdrasil: Subscription and publishing server for Elixir applications.
but it’s not maintained anymore
giddie
I wasn’t aware of yggdrasil. Thanks
My main concern would be that it requires subscription. If a process relies on durable messaging, but it dies, what happens to the messages that enter the queue while the process is down? I would like guarantees that the messages will be delivered to the process once it’s up again. And it looks like that kind of guarantee may not be available for yggdrasil.
a3kov
What’s wrong with Oban for this usecase ?
RabbitMQ+Broadway is a great combo, but it adds complexity to the deployment
You don’t need Broadway for this, Broadway is more like a job queue.
giddie
There’s quite a bit of boilerplate around each job type in Oban. You need to define a module for each kind of job you want to process. Additionally, I ran into frustrations with Oban due to the way jobs are locked for processing. In my opinion, it should not be necessary for the Lifeline plugin to exist. A failed node should cause a job to return to the queue quickly and automatically.
I’m not sure I follow. You’re right that Broadway isn’t strictly necessary, especially if you don’t need concurrency. It does have some interesting features, though.
Messages are ordered by primary key and deleted from the queue as they are processed. So messages should not be missed, but there are no guarantees about messages being processed in strict id order.
But assuming that messages are inserted serially, the messages are then guaranteed to be processed in the order they were inserted, so long as the processor concurrency is 1 (which is the default). In this case a processing error will also block the queue, i.e. the processor will not skip to the next message, to ensure message order is preserved.
If processor concurrency is configured > 1, then it will infer that it’s fine to reorder messages, and failed messages will not block the queue. In this case, if a backoff function is configured, it will apply to individual messages (instead of the whole queue), and instead of blocking the queue they will be re-enqueued with metadata tracking number of attempts, and a “process_after” timestamp.
QueueProcessor module documentation outlines the available configuration options.







