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)

Reply via email to