Building a distributed log using S3 (under 150 lines of Go)
avi.im
avi.im
I wonder, is there a formal definition for a set of primitives that allow you to build an ACID database? Assume an API of some kind (in this case, S3) that you can interact with - and provides, I don't know, locks, % durability, etc.
What would make you say, 'Having those primitives, I CAN build an ACID database on top of it'?
Sirupsen goes through the basics in [2] (video)
I did a talk about building on object storage with lots of examples [3] (video) [4] (slides with links)
[1] https://engineering.linkedin.com/distributed-systems/log-wha... [2] https://www.youtube.com/watch?v=RFmajOeUKnE [3] https://www.youtube.com/watch?v=ei0wwTy6_G4 [4] https://static.sched.com/hosted_files/kccncna2024/c8/Object%...
Depending on how you structure the underlying pages, you’ll get to decide how availability at the log level translates to availability in your user/app-facing interface and whether you will end up sacrificing consistency, availability, or partition tolerance.
Basically, S3 with its recent consistency guarantees and all-new CAS support is sufffiicent in-and-of itself. But for anything other than the most basic (least amount of data, lowest frequency writes, etc) you’ll need a considerable amount of magic to make it useable.
The most straightforward approach would be to use the existing whole of another database but swap out the backend and then tweak the frontend accordingly. SQLite lets you use custom vfs providers (already used to provide fairly efficient SQLite over http without serving the entirety of the database, but previously not for writes) and with Postgres you can use foreign data wrappers. But in both cases you’ll basically have to take out a lock to support writes, either on a page or a row (either risk lots of contention or introduce a ton of locking and network latency overhead).
There are plenty of ways to scale out traditional RDBMS, but serverless offerings make it so easy to scale out.
Is a table format which can be use via trino, spark, flink, java apis, pything API?
GCS, unfortunately, does not support copying a range. OTOH, it has long supported object append through composition.
The challenge with both offerings is that writes to a single object, and writes clustered around a prefix, are seriously rate limited, and consistency properties mostly apply to single objects.
GCS 32 objects, and Azure blob storage, afair 5k objects. Both can do an operation similar to what you described for S3 with various alternatives of read at offset + length and rolling those up.
In all cases, you end up always rolling up into a new key that isn’t available for read until roll up is done. It’s kinda useless for heavy write scenario.
Compare that to your normal fs operation. Write at offset to an existing file with size smaller that offset will just truncate the file to the offset, and continue writing.
You create a multi-part with 3 chunks: the 1st part is a copy range of the prefix, the 2nd part the bit you want to change, and the 3rd a copy range of the suffix?
And yes, all of this is useless for heavy (and esp. concurrent) writes.
GCS compose can also have target be one of the source objects, so you can append (and/or prepend) “in place.”
For GCS compose the suffix/prefix need to be separate visible objects (though you can put a lifecycle on them). For multipart, the parts are not really objects ever.
The performance isn't great because because updating the “index” is slow and rate limited, not because the APIs aren't there.
Downside is the 250ms latency. But then again, a fair amount of workloads can deal with 250ms of latency.
tmp/uploads/object-0, tmp/uploads/object-1024, tmp/uploads/object-2048
would be a rolled up object of size 2048MB + whatever is in the object-2048 file.
definitely!
I plan to add a batch write API. Also, an API where it buffers till it reaches certain size or a timeout to write to S3
tracking the batch write here: https://github.com/avinassh/s3-log/issues/3
(The goofys is faster than s3fs because it’s not totally POSIX compliant.)
I am a big fan of setting up a full stack on a COMMODITY SERVER and it just working. You can outsource your video transcoding and storage to vimeo, AWS etc. but you don’t have to! You can use your own hard drive in your home to run a social network for instance.
Or is it best effort?
Currently, there is durability guarantee. The call returns only after successful write to S3. This is close to fsync in a single node system. I plan to add a relaxed write mode
I am tracking the issue here: https://github.com/avinassh/s3-log/issues/4