139 karma · joined October 24, 2017
The page you linked to talks about how HLC is used to get distributed transactions to work, and there are differences here between the different databases (Spanner, CRDB, YugabyteDB, etc).
Disclosure: co-founder/CTO of YugabyteDB project
Actually, the issue is about the max clock skew guarantees (as opposed to the average or median). Even a single violation of this breaks the ACID semantics. So, we do use chrony, but need all this to ensure there is a max guarantee. We would totally have adopted an existing solution - we did look at all alternatives available.
Disclosure: one of the founders of the YugabyteDB project
We think a lot about this exact question. Here are some of the things YugabyteDB can do as a "modern database" that a PostgreSQL/MySQL cannot (or will struggle to):
* High availability with resilience / zero data loss on failures and upgrades. This is because of the inherent architecture, whereas with traditional leader-follower replication you could lose data and with solutions like Patroni, you can lose availability / optimal utilization of the cluster resources.
* Scaling the database. This includes scaling transactions, connections and data sets *without* complicating the app (like having to read from replicas some times and from the primary other times depending on the query). Scaling connections is also important for lambdas/functions style apps in the cloud, as they could all try to connect to the DB in a short burst.
* Replicating data across regions. Use cases like geo-partitioning, multi-region sync replication to tolerate a region failure without compromising ACID properties. Some folks think this is far fetched - its not. Examples: the recent fire on an OVH datacenter and the Texas snowstorm both caused regional outages.
* Built-in async replication. Typically, async replication of data is "external" to DBs like PG and MySQL. In YugabyteDB, since replication is a first-class feature, it is supported out of the box.
* Follower reads / reads from nearest region with programmatic bounds. So read stale data for a particular query from the local region if the data is no more than x seconds old.
* We recently enhanced the JDBC driver to be cluster aware, eliminating the need to maintain an external load balancer because each node of the cluster is "aware" of the other nodes at all times - including node failures / add / remove / etc.
* Finally, we give users control over how data is distributed across nodes - for example, do you want to preserve ASC/DESC ordering of the PKs or use a HASH based distribution of data.
There are a few others, but this should give an idea.
(Disclosure: I am the cto/co-founder of Yugabyte)
(I am the CTO of Yugabyte) Your points are all completely valid. Just wanted to add my 2 cents.
With YugabyteDB specifically, we are more than just PostgreSQL wire-compatible, we "reuse" the upper half of PostgreSQL to support almost all PG features (examples: stored procedures, triggers, partial functions, etc). So the aim is to build something that has "almost all PG features" while being able to "run cloud-native - with HA, scale and geo-distribution" - our hope is that this allows YugabyteDB to really become a viable option instead of PostgreSQL when apps are being built for the cloud. Here is a blog post on the benefits we realized reusing PostgreSQL: https://blog.yugabyte.com/why-we-built-yugabytedb-by-reusing...
Thanks for your feedback! Not sure when you tried yugabyteDB, but our serializable isolation level and YSQL API (which is needed to exercise serializability) were in beta till a couple of days ago. That said, if you can share some feedback, that would help us out immensely. All kinds of feedback welcome - be it about the product or why you feel we are not transparent. Absolute transparency has always been our goal, your feedback will definitely help us improve.
(cto/co-founder)
There are two benchmark workloads, simple inserts and secondary index. The simple inserts workload 50M unique key-values into the database using prepare-bind INSERT statements with 256 writer threads running in parallel. There were no reads against the database during this period. The secondary index inserts does the same things against a table which has an index on the value column (forcing each operation to transactionally update the primary and index tables). We use this simple workload frequently to test scalability of YugabyteDB. If interested, here is the sample app repo: https://github.com/YugaByte/yb-sample-apps
Here is a previous version of the benchmark comparing YCQL (no YSQL here, it was not GA then), it has more details about the scenario we are going for (this part applies to this benchmark also): https://blog.yugabyte.com/yugabyte-db-vs-cockroachdb-perform... The above post goes through more workloads than this one - we have not had the time to run everything on YSQL yet. The aim of this post was mainly to explore write perf and scalability vs Amazon Aurora.
> As another note, I believe Aurora added multi-master support recently. How does that compare as well?
Good question. Multi-master deployment sacrifices consistency with last writer wins semantics, so we did not benchmark against that. For example, the benchmark driver can write to one node and read from another before the data was replicated. But may still be interesting from a tradeoff perspective (perf vs consistency).
Also, note that we have also just announced multi-master support between separate YugabyteDB cluster - so a benchmark is probably something we should do at some point anyway!
* The node given to the JDBC driver could be down, or nodes could get added/removed over time. This makes "discovery of nodes" a problem.
* We could use a random round robin strategy to connect to nodes - where client connects to a random cluster node which internally connects to the appropriate node. This would result in an increased latency, and also an increase in the net number of connections needed.
These may not matter if the load balancer is smart like in the case of Kubernetes (and where the extra hop becomes mandatory as well). But for non-k8s deployments, these help.
Peeling the onion one more layer, there are two underlying features that are required:
1) Move to a threaded model to be able to scale instantaneously. This is already the case with YCQL, we're planning on doing this for YSQL.
2) There needs to be a change to the client drivers in order to become smarter to deal with how the app connects to the DB. In fact, we're working on enabling this for the Spring (Java) ecosystem by enhancing the JDSB driver (still in the early stages): https://github.com/YugaByte/ybjdbc
(cto/co-founder)
* We implemented a feature to share memory between these connection handling processes so that makes it a bit more efficient
* Longer term we are thinking of switching to a thread based model. This is how our other API YCQL works.
(cto / co-founder)
Thanks for your comment... had a couple of clarifications. Our take on the open-source licensing of CockroachDB/MongoDB has no implication on the features of these products. So in essence, you raised two separate points (one about the licensing, the other about the dashboard/features), happy to address both.
> Question about the dashboard
The portion of the dashboard that is not related to the orchestration is going to be open-source. This is a work in progress (it needs engineering work to separate these, and we're almost there). In fact, would love to have you beta test this if you're interested - please join our community slack and holler at me.
Note that we already had the core-DB open source - now we are making previously enterprise-only features open (like encryption at rest/on wire, distributed backups and read replicas). For comparison, CockroachDB (and MongoDB) had enterprise features - which were closed (and remain so). They made changes to the core-DB as well to prevent competition, effectively making it non-open-source (if you go by the definition of open-source).
Note that we're keeping "managed service" portion called YugaByte Platform under a closed license - this is the part that "manages" the cloud experience like creating nodes automatically, configuring security groups, etc. Not sure if these other companies have an equivalent product to this.
>> However, the reality is that CockroachDB is not yet at the levels of adoption where AWS would be interested. So why did Cockroach Labs make the change?
> Because it's too late once they get to that stage. The licence needs to be changed ahead of time. Because companies generally try to plan ahead more than a few months into the future.
> I was following along with you until this point, but now I feel like you're trying to hard to paint these companies in a negative light and not give any benefit of the doubt.
We named changes done by 4 companies (Elastic, Confluent, MongoDB, CockroachDB). Two of these are positive changes (Elastic and Confluent made only enterprise features closed, while the core is still open). MongoDB and CockroachDB closed their core as well with restrictive licenses. So we're simply calling that out.
MongoDB is being offered as a service by Azure CosmosDB and AWS. The licensing change done by MongoDB did not deter Amazon. Similarly, the enterprise features in Elastic were rebuilt by AWS and open-sourced into a fork. So the lesson here is - if a cloud provider wants to, they will build it anyway. Additionally, in the case of CockroachDB and YugaByte DB - AWS has Aurora, which is over a billion dollars in yearly revenue offering the same API as Postgres, so no net new functionality here. The few extra features offered by both these products can easily be replicated by AWS given the above. AWS is in fact more likely to build on Aurora rather than taking any of these services. So in the light of this, the above point of view seems fair.
PS: Will share internally about the link not working, thanks for pointing out!
Tuning to perform these efficiently in a distributed manner is a work in progress (this is wrt the time complexity question). The hints used by a traditional RDBMS would not be effective in the distributed case - optimizing this requires changing the optimizer and doing more "push-downs" to optimize the queries. Currently, depending on the query, these may or may not be fully optimized - but we hope to have good coverage soon!
This might help answer some of your questions: https://docs.yugabyte.com/latest/introduction/#what-are-the-...
Please look at the trade-offs in the above section (vs SQL, vs transactional NoSQL, vs eventually consistent NoSQL).
(founder/cto at YugaByte)
True, but note that the comparison only focuses on SQL (as it related to PostgreSQL) features and not any DB-specific features.
The YugaByte DB specific pieces are not included as well - for example, support for YCQL (Cassandra-compatible) and YEDIS (Redis-compatible) APIs to the DB.
(disclosure: founder/cto at yugabyte)
This is a very insightful suggestion, thanks for raising that! We had considered many of these variants until finally, we concluded that fully open is the best way.
PostgreSQL (which is the database on fire right now) got to this spot by being fully open and permissive - and embracing all forms of competition. In fact, PostgreSQL got rewritten from Lisp to C (which begs an interesting question - what is a database? The code or the query language? Anyway I digress).
We felt if we want to build something as foundational as PostgreSQL for the cloud, then we need to be as open.
I am the CTO/Founder of YugaByte and author of the above post. Thanks for your comments, glad you liked the post! You make a great point about YugaByte DB using Apache Kudu to start out, but wanted to clarify a few things.
* This is a post about our architectural decisions in building a distributed SQL database, and less so about how we actually implement it. Hence there is no mention of any starting points as far as the codebase.
* Also, we have used the RocksDB and PostgreSQL codebases in their entirety in addition to Apache Kudu, and make no secret of this fact (https://docs.yugabyte.com/latest/architecture/concepts/ackno...).
As @aphyr had mentioned, any NTP-alike system would work. We can update the docs to mention PTP, we do work with AWS Time Sync as well (which uses Chrony).
PG is planning the pluggable storage API for version 12.0, but these will only address user tables. This is certainly essential, but not sufficient for true distributed SQL. Some of the other pieces are: * Pluggable storage for system tables * Ability to create the initial set of system tables in a pluggable manner (initdb equivalent) * A number of other changes to the upper half of PG in order to make it "stateless" and scale-out, such as modifications to how certain operations are executed, which locks are held, what operations are pushed down to the underlying storage layer, etc. * Support for different types of indexes * Enhancing the optimizer to understand different storage types (in this case a distributed store)
In our current implementation, we have tried to make APIs for the above where ever possible. We need to explore how we can make runtime hooks for the other places. The plan is to see how to contribute these to the open source over the longer time frame so that we can create a self contained extension (though we are nowhere close to that today). Since our PG modifications code-base is open-sourced, the hope is that it would make contributing these changes back to PG easier.
From the bailis.org link you posted:
> Linearizability is a guarantee about single operations on single objects.
The place we references "linearizability" (in the Jepsen blog and docs) are in the context of single row-key operations; I can update our docs to clarify that further.
For multi-key transactions, our docs clearly point out that we support Snapshot Isolation now and Serializable Isolation is in the roadmap.
From the post by Daniel: << CockroachDB, to its credit, has acknowledged that by only incorporating Spanner’s software innovations, the system cannot guarantee CAP consistency (which, as mentioned above, is linearizability).
YugaByte, however, continues to claim a guarantee of consistency. I would advise people not to trust this claim. YugaByte, by virtue of its Spanner roots, will run into consistency violations when the local clock on a server suddenly jumps beyond the skew uncertainty window. >>
The statement about YugaByte DB is incorrect.
1. With respect to CAP, both Cockroach DB (https://www.cockroachlabs.com/blog/limits-of-the-cap-theorem...) and YugaByte DB (https://docs.yugabyte.com/latest/faq/architecture/#how-can-y...) are CP databases with HA and there is really no difference in the claims.
2. With respect to Isolation level in ACID, YugaByte DB does not make the linearizability (called external consistency by Google Spanner) claim. YugaByte DB offers Snapshot Isolation (detects write-write conflicts) today and Serializable isolation (detect read-write and write-write conflicts) is in the roadmap (https://docs.yugabyte.com/latest/architecture/transactions/i...).
3. We have publicly claimed that we do rely on NTP and max clock skew bounds to guarantee consistency. For example, slide 43 of our NorCal DB Day talk (https://www.slideshare.net/YugaByte/yugabyte-db-architecture...) we mention we are “relying on bounded clock sync (NTP, AWS Time Sync, etc).”
The point we were trying to make was that we were serving queries with low latencies from disk in a cloud environment, where folks have told us the random read latencies are higher (presumably due to multi-tenancy and virtualization).
We had run this experiment at the request of a customer who had benchmarked a number of systems and shared this concern about higher latencies especially when having to read from disk. They were very happy to see this result (pleasantly surprised even), because performance in cloud environments seems to be an issue for multiple people. So sharing widely.
BTW, we did some performance analysis/benchmarks against Cassandra here: https://blog.yugabyte.com/building-a-strongly-consistent-cas...