Ivan Andika created RATIS-2597:
----------------------------------
Summary: Local Sharding for StateMachine in a Ratis group
Key: RATIS-2597
URL: https://issues.apache.org/jira/browse/RATIS-2597
Project: Ratis
Issue Type: New Feature
Reporter: Ivan Andika
Assignee: Ivan Andika
Simply an idea.
Currently, each Ratis group contains a single Raft log and StateMachine (+
updater). This makes every transactions to be serialized. However, in the case
of key value store there are possible cases where there are some independent
shard (e.g. prefix) where operations is guaranteed to be within that shard). In
this case, the single StateMachine can cause operations across independent
shards to content with each other. In the context of Apache Ozone, currently
metadata operations are independent for different buckets, but since Ozone uses
a single Ratis group for the namespace, one bucket operations can block other
bucket operations.
We can technically create multiple Ratis groups where each Ratis group is a
shard. This is usually done by either creating a static sharding mechanism
(e.g. divide the namespace to 256 shards) or using a separate shard manager
that will be able to do split, merge a single namespace to multiple shards
(e.g. TiKV, HBase regions, CockroachDB, etc). However, usually this causes more
complexity (especially the later), since we have to manage N Raft groups, each
with their own leader and follower combination. Additionally, in the case of
Ozone, we want OM to still share a single RocksDB state so that we can halt the
Raft group if there is a InstallSnapshot. This is very difficult if we have
multiple independent Raft groups.
What if we can make a single Raft group be transparently contain multiple Raft
logs, each with a single StateMachineUpdater. The Raft group contains only a
single leader and followers combination, but underneath there are sharded logs
that can be used to improve the paralellism. This way, the Raft group is
transparent to client and client does not need to keep track of different Raft
groups with different leaders. One requirement is that each request now needs
to specify a shard key so that it can be routed to the correct StateMachine
shard.
Drawbacks
* The leader for sharded Ratis is shared among the shards which can be the new
bottleneck
* Raft implementation can become more complicated since we need to keep track
of the states of each shard.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)