Reliably replicating data between Postgres and ClickHouse
benjaminwootton.com
benjaminwootton.com
It works well. Their team is great. I feel a bit spoiled having had as much access to the engineering team during the private beta as we've experienced.
It's great for use cases where it makes sense to sync postgres tables across to clickhouse without denormalizing them. PeerDB can transform rows in a single table sent via CDC using a lua scripting language, but it can't (yet!) denormalize data into clickhouse that is stored in 3NF on Postgres across multiple tables.
On the clickhouse query side, we end up wanting denormalized data for query performance and to avoid JOINs. It's frequently not a great idea to query in clickhouse using the same table structure as you're using in your transactional db.
In our experience we sync a few tables with PeerDB but mostly end up using app-level custom code to sync denormalized data into Clickhouse for our core use-cases. Most of the PeerDB sync'd tables end up as Clickhouse Dictionaries which we then use in our queries.
PeerDB works well and I like it for what it is. Just don't expect to be satisfied with querying in Clickhouse against the same table structure as you've got in Postgres unless your data size is tiny.
Curious to know about how others are using it and the architectures you've developed.
Overall, what you shared makes sense for use cases like yours. However, there are other scenarios—such as multi-tenant SaaS analytics running large-scale workloads with PeerDB/PostgreSQL CDC. In these cases there are 100s of tables across different schemas that are synced using CDC. Some customers denormalize tables using materialized views (MVs), which is a powerful feature in ClickHouse, while others power dashboards directly with JOINs using the recent JOIN improvements in ClickHouse and suitable/optimized order keys (tenant_id,id).
When dealing with 100s to 1000s of tables and a heavily relational schema, building dual-write pipelines with denormalization becomes extremely difficult—especially when the workload involves UPDATEs.
We have many customers falling in the above bucket, replicating multiple petabytes of data to ClickHouse. A few customer deep dives on this are coming soon! :)
Side note: We are tracking support for in-transit transformations as a future feature. However, MVs are the way to go—more of an ELT approach.
> In our experience we sync a few tables with PeerDB but mostly end up using app-level custom code to sync denormalized data into Clickhouse for our core use-cases.
Have you explored dbt? You may find that using custom code is not scalable, and that dbt solves this exact problem.
dbt is as I understand it for batch processing transformations on a set schedule.
In our setup, we use app ingestion to send all the denormalised data into Clickhouse using async inserts and Debezium/Kafka/Kafka engine and materialized views to sync a few Postgres tables into Clickhouse. 2 of the replicated tables are in the order of billions of rows, and are used in 20% of the queries (usually directly and less frequently with no more than 1-2 joins). Everything else queries the denormalised tables (usually no joins there, only some dictionary usage). Overall query performance is great, although it would have been even better since we use replacing merge trees and final.
The 2 main issues that we are facing are:
- we need to periodically cleanup the deleted rows from the replacing merge trees, since the application does lots of upserts and deleted rows just stay there.
- there is not much flexibility in the ordering keys of the replicated Postgres tables, unless you enable full replica identity. We took that performance hit (although nothing really noticeable in Postgres side) in order to have some flexibility and better query performance in the replicated tables in Clickhouse.
1. For deleted rows, you can create policies to simplify querying. However, periodic deletions are still necessary. We've been optimizing lightweight deletes/updates to improve performance, which should help with automatic deletions.
2. For the second issue, refreshable materialized views with different order keys than raw tables are an option worth considering. However, having it in real time for tables with billions of rows might not be viable. That said, processing within tens of minutes to a few hours could work. We're tracking that the order key serves a dual role—as both a deduplication key and a skip index—which is the root cause of this issue of enabling REPLICA IDENTITY on Postgres side.
Separately, working on a guide covering best practices for Postgres to ClickHouse data modeling, detailing these concepts further. More on this coming soon!
This is exactly the kind of problem we've been solving with a few of our customers. With Timeplus, we can listen to Kafka and then do streaming joins to create denormalized records to send downstream to ClickHouse. Traditionally we did this with stream processing and this would build up really large join state in memory for when cardinality on the join keys would get very large (think 100s of millions of keys).
This has recently been improved with two additional enhancements: 1. You can setup the join states to use hybrid memory/disk based hash tables in Timeplus Enterprise if you still want to keep the join happening locally (assume all data in the join is still coming in via Kafka) and maintaining high throughput
2. Alternatively, where you have slow changing data on the right hand side(s), we can use a Kafka topic on the left hand side and do direct lookups against MySQL/Postgres/etc on each change on the LHS. This takes a hit throughput but may be ok for say 100s of records per second per join. There's an additional caching capability with TTL here to allow for the most frequently accessed reference data to be kept locally so that future joins are faster.
On additional benefit from using Timeplus to send data downstream to ClickHouse is being able to batch appropriately so that it is not emitting lots of small writes to ClickHouse.
Don't be confused by the timeseries branding.
the data stays in PGDB - TSDB is an extension installed onto the data base server
Of course if you come to our cloud you're going to have to do some sort of migration effort but that shouldn't be more complicated than going from one postgresdb to another.
I'm all for keeping as much as possible in your initial Postgres deployment as possible. If your team isn't having to work around things and things "just work" it's a wonderful conjunction of requirements and opportunity. It's incredible how much you can get out of a single instance, really remarkable. I'd also add it's still worth it even if there is a little pain.
But I've found that once I cross about 8-12 terabytes of data I need to specialize, and that a pure columnar solution like ClickHouse really begins to shine even compared to hybrid solutions given the amortized cost of most analytic workloads. This difference quickly adds up and I think at that scale really makes a difference to the developer experience that a switch is worth the consideration. Otherwise stick to Postgres and save your org some money and more importantly sanity.
You reach a point when you have enough queries doing enough work that the extra I/O and memory required by PAX/hybrid becomes noticeably more costly than pure columnar, at least for the workloads that I have experience with.
ClickHouse is now in my toolbox right alongside Postgres with things to deploy that I can trust to get the job done.
- you'd probably at least want a read replica so you're not running queries on your primary db
- if you're going to the trouble of setting up a column store, it seems likely you're wanting to integrate other data sources so need some ETL regardless
- usually column store is more olap with lower memory and fast disks whereas operational is oltp with more memory and ideally less disk io usage
I suppose you could get some middle ground with PG logical rep if you're mainly integrating PG data sources
https://www.timescale.com/blog/how-we-scaled-postgresql-to-3...
(Post is a year old, IIRC the database is over one petabyte now)
In general the argument I was originally trying to make is not to never use ClickHouse, I think it's a great product. But if you already are on postgres, it might just be easier to give Timescale a try than to adapt everything to work with ClickHouse right away. There is more to consider here than raw query speed.
And while I'm sure the systems behave differently scale and speed wise, I also wouldn't say Timescale looses straight up, there is situations where Timescale is faster and if it really breaks down for a use-case nothing stops you from still doing the postgres to ClickHouse migration. In the end timescale is just a better postgres, so there is no lock in.
(Disclaimer: This is Sai from ClickHouse/PeerDB team)
- Citus is AGPLv3 https://github.com/citusdata/citus/blob/v13.0.1/LICENSE
- Hydra is Apache 2 https://github.com/hydradatabase/columnar/blob/v1.1.2/LICENS...
- Timescale is mostly Apache 2 https://github.com/timescale/timescaledb/blob/2.18.2/LICENSE
As far as I understand Hyda + DuckDB it is a higher level add-on onto postgres than timescale is. This is on the one hand nice since you can just put it into an existing database without any migration effort whatsoever. But this also means they likely don't really interact with systems like the storage engine. For example Timescales deeper integration allows us to actually store the data differently on disk which allows for saving space via compression.
But it is much less performant than CH.
Always choose the right tool for the right job, Clickhouse is amazing software. I just wanted to mention that if someone currently runs analytics queries via postgres and runs into performance issues, trying out timescale doesn't really hurt and might be a simpler solution than migrating to Clickhouse.
Looks like you've expanded into vector indexing - https://github.com/timescale/pgvectorscale - and an extension which bakes RAG patterns (including running prompts from SQL queries) into PostgreSQL: https://github.com/timescale/pgai
And yes you are correct, pgvectorscale scales pgvector for embeddings, and pgai includes dev experience niceties for AI (eg automatic embedding management).
Would love to hear any suggestions on how we could make this less confusing. :-)
I guess that's why we have marketing teams!
Ajay already commented, that he's open to new ideas on how to frame timescale. I for myself always thought of postgres as the best jack of all trades database. It's basically the best db if you don't 100% know what's the best choice yet. Timescale expands on that and enhances postgres's capabilities even further so that use-cases which would usually call for a second storage option (e.g. analytics, vector) end up working great with just postgres itself.
I'd personally love if we also had a full-text offering akin to paradedb/pg_search so noone would ever need to host elasticsearch again. But it also doesn't make sense to spread the valuable postgres expert resources too thin.
But actively trying to simplify and remove as many gears as possible.
The parquet file is a columnar friendly friendly that can then be simply inserted to clickhouse or duckdb or even queried directly.
This script and a cron job are enough for my (not very complex) needs on replicating my postgres data on clickhouse for fast queries.
https://clickhouse.com/docs/sql-reference/table-functions/po...
https://clickhouse.com/docs/materialized-view/refreshable-ma...
Additionally, you can set up incremental import with https://clickhouse.com/blog/postgres-cdc-connector-clickpipe...
if you use AWS, you can upload to s3 and query via Athena. AWS Glue is a little klunky but it works and if your query load is small then it's cheap and very reliable.
I'm curious if you have data that backs this up, or if it's more of a "gut feeling" sort of thing. At first blush, I agree with you, but at the same time, by doing it at the application level, it opens up so many more possibilities, such as writing "pre-coalesced" data to the data warehouse or pre-enriching the data that goes into the data warehouse.
Secondly, OLAP/DWH systems aren’t as friendly as OLTP databases when it comes to UPDATEs/DELETEs. You can’t just perform point UPDATEs or DELETEs and call it a day. So why not let a replication tool handle this for you in the most efficient way.
With PG17+, this shouldn’t be an issue due to failover slots.
Another important observation: Aurora Postgres, which is used by many customers, persists slots during failover, so this isn’t a problem at all. PeerDB has built-in retries that resume from the last committed source LSN.
(1) https://debezium.io/documentation/reference/stable/connector...