When an SQL Database Makes a Great Pub/Sub
threedots.tech
threedots.tech
It’s curious that over those projections, we then build event stores for CQRS/ES systems with their own projections mediated by application code.
Let’s also mention the journaled filesystem on which the database logs reside. And the log structure that your SSD is using internally to balance writes.
It’s been a long time since we wrote an application event stream linearly straight to media, and although I appreciate the separate concerns that each of these layers addresses, I’d probably struggle to justify them from first principles to even a slightly more Socratic version of myself.
Also, to my knowledge, the logs in a DB are not kept forever. Instead they are trimmed as soon as reasonable. It starts to smell a little bit like a https://en.wikipedia.org/wiki/Log-structured_merge-tree
That's really logical. From the view of the application there is no transaction log, only a table. It's an implementation detail of the database.
The application wants similar guarantees a log can provide, so they build their own.
[1] https://www.confluent.io/blog/turning-the-database-inside-ou...
Now, I wholeheartedly agree that Kafka is more scalable, but I think the key point here is that there is no particular law of nature as to why that is the case. It may just be an historical accident of how RDBMS - and PosgreSQL in particular - have evolved. Further: many of the properties of Kafka are in fact also desirable properties for the PostgreSQL transaction log.
My take on the inopinatus observation, together with the Samza article [1] mentioned on this discussion, are as follows. You can think of Postgres as two "products" (bounded contexts, if you like):
- a stream-based, possibly replicated, transaction log;
- a projection of that transaction log into relational calculus, plus all of the associated machinery.
Thus far we never had the need to think of these as clearly separate "products", but Kafka makes it obvious that they are. In truth, the amount of tools processing WAL outside of PostgreSQL were already hinting in this direction; Kafka just made it obvious.
From this perspective, it seems a tad expensive to take the original transaction log, convert it to a RDBMS representation, then convert it to events and, in some cases, then store it as an event stream in Kafka. It would be much more efficient to simply use the original transaction log directly - and this is why, to me, even Debezium [2] / Bottled Water [3] appear to be one layer too many. To the best of my understanding, this line of reasoning is also line with the observations in the Samza article [1]. Where I believe I differ from the article is in thinking that the RDBMS representation also adds a lot of value to applications - I see both having a role (e.g. streaming vs batch processing sort of thing). I think this would derail the present discussion too much, so I won't go in to it.
In conclusion: to the untrained eye, it seems that the right thing to do is to extract the transaction log out of PostgreSQL and make it as scalable as Kafka. Then, allow for it to log "things" which are not necessarily "projectable" into the relational plane. PostgreSQL then becomes just a client of the transaction log, together with other "kinds" of clients. I suspect that this is what will ultimately happen, but the engineering work required will probably span a decade or more.
My 2 Angolan Kwanzas, at any rate.
[1] https://www.confluent.io/blog/turning-the-database-inside-ou...
Kafka is also just more optimized for what it does. Postgres is a superset of what kafka does, so kafka is unsurprisingly better able to optimize for its usecase. It has a zero-copy protocol that can shuttle data to/from disk to/from the network without bringing it into memory (using the sendfile syscall). It doesn't wait for disk writes when doing writes, because it achieves durability via replication.
Also, don't forget the things you'd have to do when implementing consumers. How will you load balance streams between consumers? Meaning, if you have a stream and you want multiple consumers to burn it down at a time, how can you make sure they aren't duplicating work and they can handle the consumer group growing/shrinking? How will you handle checkpointing where each consumer tracks what they've done so far? What about streams' data rolling off?
All of this is doable with pg, but you'd have to implement it yourself. With kafka and its client drivers, this is handled for you.
I would update all references of the former to the latter.
The academic/enterprise database space has been discussing and tackling the types of questions you raise for decades. I don't think that is a useful lens to evaluate this article which is effectively a "tips on when to use our GoLang SQL Pub/Sub layer".
One must be careful not to use this as an excuse, of course, and keep an eye on the scaling concerns and certain other details. SQL-as-pubsub has certain well-known issues and anyone using it this way ought to be aware of them. But it's a thought worth having.
I've got a system I'm managing where Cassandra is the backend. I've got about 5 "documents" (in the MongoDB sense, let's say) I want to store in the system. I don't put up a whole "document DB" for them, I just have a table in Cassandra. I have in my entire system, one distributed lock I'd like to have per certain resource, of which I expect there to be single digit numbers of that resource over the lifetime of the system. Cassandra is not a great distributed locker, but it does work (and as far as I can tell, done properly, is also correct), so rather than install an entire distributed lock server, I use Cassandra.
Should this ever turn out to become a mistake, all code that uses either of these functions is cleanly isolated and I can easily swap them out later. I am aware of the possibility that could happen in the future and have prepared for it. In the meantime, I've avoided two entire systems being poorly deployed and understood in favor of the one system that is well-deployed and understood by the team.
Such as? Other than performance.
I have an application at work I have to use (not something I wrote) based on an old version of Postgres that was experiencing this until quite recently when they upgraded, and I'm still not 100% sure that it's fixed because it hasn't been enough time. (I do expect it to be fixed, though.)
The aforementioned Cassandra is pretty much catastrophic in this use case, from what I've seen; I didn't try to use it as a message queue, all I did was write some unit testing that involved adding some test records during the test and removing them, then running a whole bunch of other tests that would add records to the same primary key and remove them, simply to clean up for the subsequent tests, a very similar pattern as using a table as an event queue from the database's point of view even though we humans see that as quite different, and I very quickly noticed my tests were running noticeably slower every execution, on scales that compared to what Cassandra can nominally handle were very, very tiny. For the testing, I found I could use TRUNCATE, which tells Cassandra "Hey, I'm not just happening to delete every row in this table, I'm actively nuking it, so please clear it out", resulting in Cassandra completely wiping the table tombstones and all, which was fine for this unit test case where they run synchronously relative to each other and I have clear points where I can say "just throw the table away now, tombstones and all", but that's completely unsuitable for an event queue use.
Somewhat ironically, you sometimes want the more "primitive" table types for this use care. IIRC there was a time when if you were using MySQL, you'd want to use a MyISAM-type table for the queue table, because it pretty much works as you expect and deleted rows are actually deleted and gone, but InnoDB used the "higher performance" tombstone-based approach, which is must faster, as long as deletes are relatively rare, but much slower if they aren't. Know your DBs!
Do databases keep the whole transaction log forever? It seems like that could keep growing forever even when the tables stay at constant size.
For your question about the logs growing forever, there are usually points in the process where a transaction log is saved off somewhere else and then can be overwritten, but on a transactional system the logs over a period of a week or two can sometimes be many times larger than the actual data stored at any given time, yes.
It might work, but it’s not the general case and you might spend more time to debug your table then to write the code to use a real queue.
And I’ve also seen people build their own queueing engine for a few hundred tasks per day. Why don’t they just choose one of the very good open source solutions?
If you already need a database for something else, using the DB as a Queue means you don't need to list {mqFlavorOfChoice} as a requirement for new hires. You also don't have to manage that extra infrastructure. Of course, you are putting additional load on the DB.
Mind you, I'm speaking of a pub-sub type queue and not a FIFO here. You can do FIFO queues in DB as well of course, it's just not as compelling of a story nowadays.
Also way easier to look at and 'poke' a Database queue if you need to. The queries are also not really difficult to write for a general purpose use case.
And how hard it is to clean up the data at this point in time depends entirely on the kind of system you're working with.
If I was building a new system that required a queue I'd definitely put in the same postgres db as the rest of the data until I had a good reason not to.
If you don’t take this into consideration then you’re detracting from the business to satisfy another need.
Polling isn't a huge issue to begin with, and is mitigated with LISTEN/NOTIFY (on certain DBs). Inserts with indexes are not a performance problem at the scale of most applications. A separate messaging service won't prevent you from building a "hugely coupled monster".
Personally, I almost always start with the database as a queue. The operational overhead of running, updating, and monitoring another entire service is non trivial. If the messaging rate exceeds the database's capabilities in the future, I'll migrate then.
(Blockchain another "oplog" that ends up caring a lot about state eventually).
It's no wonder you can use them interchangeably in many common base cases.
Messages will never arrive, arrive out of order and I don't remember the third one right now (messages will arrive late?)
Some messages might still be too late and get missed, but most of them should get through, which is what every messaging service is currently doing
The application logic can indeed solve it, "eventually consistent systems" are a thing. My main goal was figuring out if this was impossible in SQL for some reason, but my understanding is just that the implementations are usually weak and do the "read-back" they need to, to recover late messages.
Tough choice, interesting nevertheless
EDIT: More context for the above process[1]
[1]https://www.eversql.com/faster-pagination-in-mysql-why-order...
What parent mean is that there may be holes in the sequence of primary keys. What you do with pagination is that you first sort the sequence, then thrown away the first N results, and finally select only the next M results.
It will work just fine.
Recommended reading: https://www.citusdata.com/blog/2016/03/30/five-ways-to-pagin...
Are you sure that's what they mean? It's not what they said. "Monotonic" means "strictly increasing" (or decreasing) e.g. 1, 2, 5, 7 is monotonic even though it has gaps. "Contiguous" means "without gaps".
When the "forward" part of "store-and-forward" is most important then Kafka is a fine solution.
However, when the "store" part - for example you want to be able to stream historical data again, or interact with the data in different ways - is most important I have recommended HBase (+ Phoenix) as a better solution in the past.
It's an Elixir server (Phoenix) that allows you to listen to changes in your database via websockets. Basically the Phoenix server listens to PostgreSQL's replication functionality, converts the byte stream into JSON, and then broadcasts over websockets. The beauty of listening to the replication functionality is that you can make changes to your database from anywhere - your api, directly in the DB, via a console etc - and you will still receive the changes via websockets.
The article suggests Postgres’ native LISTEN/NOTIFY functionality. I tried that originally and found that NOTIFY payloads have a limit of 8000 bytes, as well a few other inconveniences.
It's still in very early stages, although I am using it in production at my company and will work on it full time starting Jan.
[1] https://www.reddit.com/r/PostgreSQL/comments/ebu6nh/message_...
An alternative implementation is provided by Debezium [1], a general solution for change data capture for MySQL, Postgres, MongoDB, SQL Server and others, based on top of Apache Kafka (but can also be used with Pulsar and others).
There's support for outbox coming as part of Debezium out of the box [2].
Disclaimer: I'm working on Debezium.
[1] https://debezium.io/ [2] https://debezium.io/documentation/reference/1.0/configuratio...
It's incredibly well written and I am using it in a project.
https://deepstream.io/tutorials/concepts/what-is-deepstream/
https://github.com/deepstreamIO/deepstream.io
(I'm not affiliated with the project.)
I'm working on a module that send notifications to a user when an alert is generated. I have PostGreSQL as the database and NodeJS is the handler and for connection pooling. Are there any good pub/sub tools that I can use. Thanks in advance.
with alert as (
select alert_id from alerts
where alert_sent_at is null
)
update alerts set
alert_sent_at = now()
from alert
returning *It works by following the MySQL binary log and triggering a reactive query based on event conditions specified by the programmer, e.g. a change in a field.