ActorDB – Distributed SQL database
github.com
github.com
However, it has a unique and initially confusing data model where the database is divided into "actors", which are self-contained shards. For example, a database modeled on Hacker News would probably have an actor per user and an actor per story, and probably an actor per thread.
Every actor acts like a self-contained database, with its own set of tables. When you want to query or update data, you first tell it which actor(s) to operate on; but unlike systems like Cassandra, the sharding is explicit, in that the shards have identifiers, and there's no automatic sharding function. Indeed, you can have actors that act like a basic database, e.g. a global, shared list of lookup values could be a single actor.
You have ACID transactions within actors, and you can also do transactions across actors, and queries can also span multiple actors. You can't do joins across actors, as far as I recall. Schema migrations also become interesting, since each actor has its own, entirely separate schema.
I gotta say... this project strikes me as very strange. If you're going to go through all the trouble of sharding your data across a set of HA actors... why not also shard behavior? Co-locate data and behavior! Isn't that a key driver Actor model? Then you wouldn't communicate with your actors using a very limited language like SQL you could send them real domain-specific messages. Most importantly, when it comes to transactions across actors... you don't need it. All you need is guaranteed message delivery to the various actor mailboxes and all that complexity melts away. You need this anyways if you want your actors to be able to actually collaborate and send each other messages. This whole project seems to be designed to eliminate many of the benefits of the actor model...
(I don't like the all-too-common HN cynicism but I also worry when I see stuff like this. Either I'm crazy and missing something obvious... or everybody else is crazy. But then much of what comes out of the Erlang space is... strange to me. That culture seems to have a very unique idea of distributed computing.)
This is a database that permits any application to read and write data without coupling it to a specific language or platform, e.g. Erlang. You just write SQL. People know SQL. The difference is that your app must be shard-aware, which can be an acceptable compromise for apps that wants to scale far. The sharding also gives you some measure of enforced encapsulation, since you can't have foreign keys (that I know) across actors.
This project might be particularly suitable for SaaS-type multitenant systems that need to have a clear boundary between each tenant.
I mean, there is some possible advantage by having all users in their own actor-namespace, for bulk deletes and concrete mapping of data-usage, but IME you're always stuck with a need to preserve the operational state of the application and the quality of its legacy data.
Regardless of sharding strategy, somehow your datawarehouse needs to be able to pull out or confirm data from the previous quarter...
Does it use snapshot isolation? Is it serializable? Is it linerizable? With all the great work Kyle Kingsbury (aka Aphyr) has done on the Jepsen tests, it's pretty clear that claiming to be "ACID" with no additional info isn't sufficient for a modern database.
-
Would you rather the software not exist at all?
I'm also curious how Raft is used and what the resulting guarantees on distributed table operations are, but I'm coming at it assuming good faith, since SQLite is... robust.
Each actor has its own WAL, and so I suppose only operations within a single actor are consistent. Multi-actor transactions use two-phase commits.
That's sort of a meaningless statement tho. Everything in the world is built on top of RAM which is serializable. The real complexity comes in how those foundational blocks are abstracted and combined.
We know that it's not performant at scale to represent an entire database log a single Raft group. So how are the Raft groups structured organized, and how and when do they cross-communicate?
[1] http://www.actordb.com/docs-howitworks.htmlhttp://www.actord...
ActorDB has no concurrency at the actor level. This makes it a poor fit for applications that have lots of concurrency around the same pieces of data. A long running read or write on a single actor will lock out any other reads or writes.
Likewise, distributed transactions lock all actors involved in the transaction until the transaction completes.
It seems like reads have to go through a round of Raft. This increases latency for reads. It also decreases throughput, although I'm not sure how big of a deal that is given the lack of concurrency.
It's unclear to me how ActorDB guarantees serializability for multi-actor transactions. You need some way to guarantee two multi-actor transactions will execute in the same order on every actor. Based on the docs, Raft is performed at the actor level and not across multiple actors. ActorDB does use two-phase commit to guarantee atomicity across multiple actors, but there's no description of how it handles serializability.
Based on my reading, ActorDB is good if you have lots of data and your queries have either low concurrency and you don't require high throughput. If you have high concurrency or require high throughput, my guess is ActorDB will be a poor fit.
If you are looking at distributed SQLite solutions there is also rqlite/dqlite and bedrock.
This sounds to be an interesting project as well.
The question to me always about "how this will makes the project that use it have little learning curve for the new recruits, easy to understand integration in the code level and low maintenance on the long run"
See:
Also, I poked around BEAM a fair bit and it seemed just fine to me. It's pretty clean and easy to understand for a fairly complex VM.
Here's a reasonable summary: https://stackoverflow.com/questions/3760881/how-is-erlang-fa...
https://ferd.ca/the-zen-of-erlang.html
It's more about assuming things will fail by default with good ways of handling it built-in at the language level. A lot of people also find its inventor's thesis enlightening. Here it is:
http://ftp.nsysu.edu.tw/FreeBSD/ports/distfiles/erlang/armst...
And for those concerned, there's also at least one project to write the native functions in a safer language to reduce their bugs a bit:
It may be a poor mans Erlang, but better than nothing ;)
Everything is built so that you can operate a telecom switch and never drop that emergency call.
Erlang (Ericson language) is not necessary, you can also program for this in Elixir, although personally I prefer Erlang.
https://github.com/tidwall/buntdb
It was very well written and easy to pick up as well.
With ActorDB you have to design your data model to shard at a high level. For example, you could shard it by user; every user would have its own database that ActorDB will shard and replicate for you. That database is for the most part separate from everything else -- it has its own tables and indexes.
ActorDB provides tools to operate on multiple actors, but they're explicit. For example, you can query across multiple actors, but this requires a small SQL declaration at the beginning of the query to select which actors to query, a bit like a "for each <actors> do <some SQL>".
You can also do transactions across multiple actors, though this uses two-phase commit (which coincidentally is the strategy used by Google Spanner), and requires some locking.
So Cockroach pretends to be a classic RDBMS (databases have tables and indexes, but most apps just use a single database per app), allowing an existing app to be ported with little effort. It would be harder to port an app to ActorDB.
Many databases are “partitioned” by user anyway so in this case the DB can be smarter if it doesn’t have to handle cross partition queries.
Seems like the idea of partitioning a MySQL database taken to the next level.
For example, it is possible for the db to allow direct access to any replica, not just the leader, and if the replica is not part of the quorum, the client will see outdated information. This may be just fine for data that seldom changes, or where a certain amount of staleness is just fine.
Stale reads technically violate the definition of Consistent in CAP. Using raft means unavailability of some subset of nodes can make the entire system unavailable, violating the definition of Availability in CAP.
So the parent comment's suggestion of serving reads from any replica would result in only a partition tolerant system, without any formal availability or consistency guarantees.
You’ll make different choices for an application that’s processing clickstream data than one that’s managing highly relational data.
2. Use eventual consistency (not suitable for every workload)
[0] https://www.microsoft.com/en-us/research/project/orleans-vir...
I have this exact use case for a new project I’m working on so I’m fine with the 1-db-per-user approach. ActorDB definitely sounds very interesting, but are there any alternatives to compare it with?
You pick your poison.
Glossing over / ignoring / implying they solved CAP is very typical for database marketing.
None of these DBs seem to be byzantine fault tolerant, though. Any examples of BFT databases?
Not sure what you mean. Obviously we can't have blockchain level of byzantine fault tolerance against malicious nodes. Only against byzantine failures. But then this is pretty much what CAP answers, where AP systems can handle byzantine failures, because they don't need consensus, while CP can't, because they do. Quorum logic is in the middle here, where you can handle some byzantine failures and have some consistency.
So not doing good query performance for OLAP/DS. And "core" DB is in Erlang.
Storage and computing are two layers in a normal application. You are free to switch computing layers as long as storage layer speaking some standard protocol.
AKKA + AKKA persistent, how can you switch the computing layer?
Deleted comment
Would like a similar solution but with Postgres on the back end instead.