Yandex open-sources its exabyte-scale big data platform
medium.com
medium.com
The most interesting part of YT is Cypress. I'm particularly interested in how they make their master cluster horizontally scalable.
However, this approach proved to be non-scalable as the memory amount and throughput of the master server soon became insufficient. To address this issue, we implemented Multicell technology. With Multicell, there are multiple RSMs called secondary masters that store information about chunks of the tables and their placement. The primary master still stores information about the distributed filesystem and transactions but is now single and non-sharded.
After a few years, the masters became overloaded again, and we implemented Portals. With Portals, one can select a subtree of Cypress and place it in one of the secondary masters. This technology is used nowadays, and home directories of some active users are hosted on secondary masters.
However, we anticipate that this approach will also become insufficient in a few years. Therefore, we are currently working on a new technology called Sequoia, which stores information about the Cypress tree shape in horizontally scalable dynamic tables.
It is hard to describe all aspects of master server internals in one comment. Therefore, feel free to join our chat at t.me/ytsaurus for further discussion!
Why not just use a database for the metadata? Something that can be sharded and has transactions like YugabyteDB/Yandex-ydb/etc?
Storing Cypress nodes in the ad hoc RSMs and information about tree in k-v storage seems a good compromise that is both scalable and allowing to implement any functionality for objects efficiently.
Having said that, this definitely does look to be an impressive feat of engineering!
I believe the engineers really spent a lot of their time building this from the ground up. I'm thankful that those large companies open-source their software and algorithms.
Both Cypress and Zookeeper are fault-tolerant distributed hierarchical filesystems that can be used for distributed coordination, but Cypress has much richer functionality.
Recall that Zookeeper's data model is just a tree consisting of homogeneous nodes that can be either ephemeral or persistent, along with a set of sessions that control the lifetime of ephemeral nodes. This simple model allows to implement multiple primitives of distributed synchronization, such as leader election, exactly-once queue processing, or two-phase commits. However, it is not always easy to integrate Zookeeper with third-party systems. For example, if you want to elect a leader via Zookeeper and use it to insert data into a database, it is mandatory that the instance remains the leader during the commit into the database, which is not easy to implement without races or some additional assumptions. In YTsaurus, transactions permeate our entire system. You can start a transaction and acquire an exclusive lock at some Cypress node (which is a way to make a leader election), and after that, the transaction becomes the leader lease. You can then modify Cypress, run MapReduce or YQL operations using the transaction as a prerequisite, lock some files and tables in the same transaction, and do many other things. Currently, we are working on the ability to use Cypress locks as prerequisites for dynamic table commits. There are many other features in Cypress that are not implemented in Zookeeper, such as symlinks, automatic expiration of unused nodes, and many others. Moreover, Cypress can be sharded using Portals about which I wrote in a previous comment, so this filesystem is scalable unlike Zookeeper. Even without sharding a single primary master of YTsaurus can hold tens of gigabytes of metadata of Cypress while Zookeeper state size is limited with hundreds of megabytes accoring to etcd vs Zookeeper comparision [1].
One major disadvantage of Cypress compared to Zookeeper is the lack of watches, so all changes tracking should be done via short polling. The good news is that Cypress is well-optimized for read queries with the possibility to read from followers and from multiple threads, so this is not a big problem. In the meantime, we are considering the possibility of adding some kind of watches to Cypress.
The big difference between Cypress and Zookeeper is the replicated state machine implementation. With all due respect to Zookeeper developers, Zookeeper was implemented over 15 years ago when the world of distributed algorithms was different. Today we see that ZAB (the consensus algorithm used in Zookeeper) has some shortcomings in failover speed and stability. There are multiple reports of Zookeeper being unstable under heavy load. In YTsaurus, we use an in-house library called Hydra for RSM implementation. This is our consensus algorithm very similar to RAFT that has proven itself to be both efficient and fault-tolerant. We use Hydra for master servers, clock servers, and tablet cells (RSMs that store data in dynamic tables). I even had an idea to implement a Zookeeper API using Hydra both to simplify migration to YTsaurus and check Hydra performance and correctness via multiple tests implemented for Zookeeper (Jepsen, for instance), but did not have enough time to finish this project.
This comment is already quite long, so I will write about the YT vs Hive comparison in another comment later on.
It's the same as when people use k8s not utilizing its full capabilities, only to be able to massively scale up when needed.
Sure it would be overkill for a lot of applications, but so is redis, react.js, etc.
P.S. I wonder if LLMs could be used to generate docs and comments for big hairy codebases. Seems that the current generation of LLMs lack context to do it, but maybe it's "just one or two more papers down the line"®...
It's truly a hard work, because we were very tightly tied to the Yandex infrastructure and we had to learn how to deploy in k8s from scratch. Also you need a new brand and clean your documentation from irrelevant things... All this takes months.
For me this seems very plausible, as for the last year they first did everything to distance from anything related to politics (e. g. they sold their news and their blogging platform to the basically state-owned VK), and then to separate Russian and overseas businesses as much as possible.
(Not suggesting this is wrong. I'm just offering an interpretation based on skimming public blog posts over the past ~2 years.)
I think they will create several spin-off open source companies (like ClickHouse Inc.) outside Russia to continue doing B2B business with the outside world.
Like https://nebius.com/about for example.
Not to mention that Yandex been operation in countries outside of Russia under different name from the beginning (like https://en.wikipedia.org/wiki/Yango_(ride_sharing))
makes sense that original engineers/founders create their own stuff via opensourcing their original work
Spoiler: because until June 22 (effectively until the sanctions hit) he was a Yandex CEO and owner of 8% shares (45% voting shares).
There is no "almighty Kremlin" that owns everything. There is, however, a set of rules you must comply to if you want to do multibillion dollar business in Russia. You either bend, or sell your business to more complacent oligarchs. Durov chose the latter, Volozh chose the first.
And so I suspect Yandex leaving it's home turf would essentially be a covert invasion of wherever they'd swarm to. Somehow, money tells me UK would be a likely target. Source: worked briefly for an exiled Russian company. Would not repeat.
> an exiled Russian company.
Either you are exhiled or you are invading, how can you be both?
Bonus "inception" points if you can make the adversary believe that you did it because they forced you to (sanctions).
I have certainly seen the kind of people you talk about, they support (or at least used to) Putin, but don't want their kids to live in Russia. I would call them more like emigrants of convenience.
It would be very interesting to see some in-depth comparisons with already-existing open source technology (like Hadoop, Hive, Iceberg, ZooKeeper) to get a sense of when and where YT could be more effective.