5 ms·
I’d love to know if anyone here has done a deep eval of Citus vs CockroachDB. They seem to be the two most promising solutions for horizontal scale-out Postgres
by hemancuso 8y ago
I’d love to know if anyone here has done a deep eval of Citus vs CockroachDB. They seem to be the two most promising solutions for horizontal scale-out Postgres and both are under very active development.
- qaq 8y ago?? outside of CockroachDB using PostgreSQL wire protocol they focus on very different use cases it seams. from CocroachDB website: "When is CockroachDB not a good choice? CockroachDB is not a good choice when very low latency reads and writes are critical; use an in-memory database instead. Also, CockroachDB is not yet suitable for: Heavy analytics / OLAP" I think Citus is actually very well suited for Heavy analytics / OLAP
- hemancuso 8y agoYes but they both target horizontial scale out and OLTP is a target for both. And yes, cockroach isn’t Postgres but it has SQL, versus something like Mongo or Cassandra.
- qaq 8y agoCitus claims OLTP but for a large cluster 2 phase commit is not a very viable approach.
- manigandham 8y agowhat do you mean?
- qaq 8y agoYou need all worker nodes to be available for 2pc to succeed. So the solution you have a standby for each node which is not very viable for a large cluster.
- manigandham 8y agoYes, the fact that Postgres is not natively distributed means that HA/replication will be very inefficient to implement. Citus is best used when transactions don't cross shard boundaries, in which case they execute on a single node and give you the low-latency to match.
- elvinyung 8y agoIt's not pure 2PC, it's 2PC on a subset of shards (layered on top of Raft). If it's true that most workloads are primary-key-based or touches a small amount of shards, it's fine.
- qaq 8y agoCould you point to info about them using RAFT? "In Citus, we looked into the 2PC algorithm built into PostgreSQL. We also developed an experimental extension called pg_paxos. We decided to use 2PC for two reasons. First, 2PC has been used in production across thousands of Postgres deployments. Second, most Citus Cloud and Enterprise deployments use streaming replication behind the covers. When a node becomes unavailable, the node’s secondary usually gets promoted within a seconds. This way, Citus can have all participating nodes be available most of the time."
- elvinyung 8y agoOops sorry, thought you meant Cockroach! The part about 2PC still holds.
- mslot 8y agoI don't think you can get around using 2PC in a distributed OLTP database, e.g. Spanner also uses 2PC for distributed transactions across shards. Fortunately, the overhead is not really that high because the prepare and commit messages are sent to all nodes in parallel. It only adds one extra network round trip. It does lower per session throughput a bit, but you can always get better throughput by creating more sessions.
- qaq 8y agoSpanner runs consensus for each key range and does not need all nodes to be available to make progress also my understanding since again it has a leader for each key range writes scale better.
- mslot 8y ago> Spanner runs consensus for each key range and does not need all nodes to be available to make progress also my understanding since again it has a leader for each key range writes scale better. Spanner uses Paxos (consensus) for replication within a key range (shard), but two-phase commit across shards: https://ai.google/research/pubs/pub39966 https://ai.google/research/pubs/pub39966 Citus relies on PostgreSQL's streaming replication, which gives higher throughput than Paxos, but Paxos has better availability characteristics. On the other hand, Paxos with leader leases as used by Spanner is similar to streaming replication both in terms of performance characteristics and short downtime during failover.
- qaq 8y agoIt's not similar in infrastructure requirements though Spanner does not have 1/2 the nodes in standby mode for HA they are actually doing work.
- deleted 8y ago[deleted]
- craigkerstiens 8y agoAt Citus we tend to focus on two use cases: 1. Analytical, but less data warehousing and more of a HTAP (hybrid transactional/analytical processing). In this case you're often ingesting a lot of data, often times sensor or log data from many endpoints, and then providing analytics across that data. The analytics needs to be up to date within minutes, and responsiveness of reports within seconds. You can see how Algolia (which powers the search for HN) uses Citus for this in their blog post - https://blog.algolia.com/building-real-time-analytics-apis/ https://blog.algolia.com/building-real-time-analytics-apis/ 2. Transactional. For a couple of years now Citus has had full transactional support when targeting a single node. Single node transactions can actually cover a breadth of use cases because it can span across tables as long as tables are co-located within the same node. We often see this is the case for multi-tenant/SaaS applications. In recent releases we also added support for distributed transactions. These transactions do have a higher overhead, but can often be hard to detangle from an existing application, thus us building support for distributed deadlock detection then adding distributed transactions. Generally we're continuing to improve and support both of those use cases and have our usage base actually pretty evenly split between the two.
- hemancuso 8y agoCan you say anything to help understand the differences between cockroach and Citus?
- manigandham 8y agoThis is a good question. We did, and they are very different databases. CRDB: Pros: natively distributed so HA and scalability are built-in, simple deployment and configuration, can run on Kubernetes for automated HA operations. Scaling tables is automatic and single-key and small-range OLTP performance is very good. Supports JSON and most data types for compatibility with most things that use Postgres. Cons: Still maturing and has bugs like `select unnest(some_array_col)` not working. Obviously cant run any Postgres extensions so SQL w/JSON is all you get. Performance on large scans is slow, they're working on this but the distributed consensus required for queries means they will never match the latency of a single-node Postgres. Advanced queries are either very slow or unsupported or every slow. CITUS: Pros: pure Postgres including extensions so you have access to advanced functionality. If you use shard key for queries, lack of distributed consensus gives low-latency performance just like single-node, but distributed transactions are still possible. Citus scales queries across all CPUs (on nodes holding the accessed data) so greatly improves query performance. Cons: only distributes data in "distributed" tables (sharded) or "reference" tables (full replicas on all nodes). All other data just sits on single master node. HA uses Postgres streaming replication, requiring an inefficient 2x increase in costs, and is not seamless with failovers. Generally requires much more maintenance because it is still Postgres. Sharding does not accept multiple columns. No columnstores so large scans can still be slow, but they have ZFS in beta. -- Summary: CRDB for simpler OLTP with very low ops overhead and great availability, scaling, and durability. Citus for advanced OLTP or OLAP, low-latency sharded access, and full access to all Postgres features.
- hemancuso 8y agoThanks! Did you consider any other options?
- manigandham 8y agoWhat scenario are you looking for? TiDB is another competitor but mysql dialect and still early, missing lots of features. For pure data-warehousing, we used MemSQL which is incredibly fast but can be expensive. SQL Server is a great all around database if you want in-memory tables, columnstores, native graph queries, full-text search, and very high performance and can live with a single-node design (with optional HA cluster).
- iKevinShah 8y agoNot fully qualified to answer this but one thing which helped me select Citus over Cockroach was Citus being natively postgres and hence extensions (like PostGIS) being supported out of the box.