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:
https://github.com/giddie/elixir-postgresql-message-queue
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. · GitHub
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.
Last Post!
giddie
Well, that’s kind of my point with this project – different usecases will have different requirements. So a library that tries to cover all the usecases is not necessarily a good idea. But some reference code can be really handy, so you can copy the bits you need and adapt it as needed.
Well no, it’s a data stream processing framework. So for processing a message queue, it’s pretty useful. It’s certainly not the only way to ingest a queue, though. A simpler GenServer for queue processing can be found here in my CQRS patterns repo, and it’s pretty much a drop-in replacement for the Broadway-based processor.
Ultimately, anything could call Messaging.process_message_queue_batch/2, so the world’s your oyster.
The main advantage of Broadway in this message queue example is concurrency – it decouples the ingestion from the processing, so one process can pull a batch of 100 messages, and distribute them across a pool of 5 workers. And of course if the work calls for it, Broadway has features such as batching and pre-fetching that could be applied very easily on top of this, since the hard work of building the Broadway Producer is already done,
Popular in Announcing
Other popular topics
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #deployment
- #library
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #channels
- #elixirconf
- #exunit
- #discussion
- #code-sync
- #javascript
- #podcasts
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #full-time-contract
- #podcasts-by-brainlid
- #ecto-query
- #elixir-ls
- #phoenix_html
- #iex
- #blog-post
- #graphql
- #genstage
- #ai
- #websockets
- #supervisor
- #elixirconf-us
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #metaprogramming
- #security
- #hex










