I think putting a stateful, persistent, transactional stores in the middle of two stateful, persistent, transactional stores is a bad idea.
The idea is to sync secondary store B to continuously be a perfect replica of A, so: A —> B. What a lot of people do, you included, is add a third store as an intermediary: A —> Q —> B. Now you have three complex pieces of software rather than two.
The thing is, you already have the state that Q covers: It's A.
For the record, we made the same mistake. We put RabbitMQ in the middle. Now we had many problems:
* What do we do if the queue loses messages? (Manually reindex from N days ago.)
* How do you know if the queue has lost messages, due to bugs, or network downtime or similar? (Well, you don't really know. If unsure, manual complete reindex.)
* What do you do when the PostgreSQL transaction completed, but it's for some reason unable to reach RabbitMQ to post a queue message? (Manually scan logs, figure out how far back to backfill, reindex from there.) Fixing this "correctly" would require some kind of two-phase commit, which queues don't support.
Lots of manual intervention required to run a system flawlessly. It's not just about running a consistent system; it's about knowing when you're consistent or not, and how to repair.
Also, there are logistical issues:
- What do you do when people run batch jobs producing millions of updates, and users want to update documents (and see their changes) at the same time? You have no recourse but to create traffic lanes — queues and more workers.
- How to run multiple queue consumers. You have to use version constraints (update only if newVersion > oldVersion), because you will end up processing updates out of order, something which only works if the original source has a version field (ES does support "external" version numbers). Turns out a queue makes these checks happen more often because you often get multiple adjacent updates for the same objects. Kafka can de-dupe, fortunately, but RabbitMQ can't.
- How do you update multiple target stores (let's say, both ES and InfluxDB)? You have to create separate queues for each of them, so that you get true fanout. Now you have added more "middlemen" on top of the first one, more stuff that can go out of sync.
My conclusion after struggling with this for a while is that the database is the truth, but it's also the only one that's coherently transactional. So one should keep the change log close to the truth, meaning in the database. You can use a PostgreSQL "unlogged" table for performance.
Secondly, you will want to use time-based polling that invokes a simple state machine that can travel back in time. Basically, you run a worker that keeps a "cursor" pointing into the database transaction log. Let it grab big batches every N seconds. The cursor doesn't need to be transactional as long as it's reasonably persistent; if you lose the cursor, your worst case is a full reindex, but if you are unable to update the latest cursor, worst is case is just a small amount of unnecessary reindexing. This worker can live side by side with the queue-based processor, and if you shard it, you can run multiple such workers concurrently.
Everything else comes out of this logic. For example: If ElasticSearch is empty, it can detect this, and set the cursor to the beginning of time. It knows how far back the cursor is, so it knows whether it's in "full reindex" mode or "incremental mode", something it can export as a metric to a dashboard. It can also backfill, by moving the cursor back a little bit. And by using a state machine you can also put it in "incremental repair" mode, where it can use a smart algorithm (Merkle trees were mentioned by someone else recently) to detect holes in the ES index that need to be filled.
Things like building a new index now becomes a trivial, because you just start a worker instance that points to a new index, but from the same truth data store; being smart about state, it will start pulling the entire source dataset into the new index. The old worker can continue indexing the old index. Once that instance is done, you can swap the new index for the old one, then delete the old worker.
Finally, to solve the problem of real-time vs. batch updates: In addition the above worker you run a separate worker that listens to a queues. Whenever a non-batch update happens (a "batch" flag needs to be indicated in all APIs and internal processes), push the ID of the affected object on a queue, but not the object itself; rather, let the worker pull the original from the store. This way, your queue (which requires RAM/disk) stays super lean and fast. Give the worker a small time-based buffer (like 1s) so that it can coalesce multiple updates if they're happening rapidly, and use an efficient query to get multiple objects at the same time. And use versioning to avoid clobbering newer data.
Of course, the system I've outlined is probably not workable for Google or Facebook, but it will scale well and will keep things in sync better than something queue-based.