IP Library › Granted Patent US 10,942,828
Granted Patent B2
US 10,942,828 · App. 16/270,048 · Granted Mar 9, 2021

Method for storing data shards, apparatus, and system

Inventors: Huaqiong Wang (Shenzhen, CN); Chao Gao (Hangzhou, CN)
Assignee: HUAWEI TECHNOLOGIES CO., LTD.
G06F11/2094G06F11/2082H04L29/08H04L67/1097
View Patent ↗
Loading inventors, assignments & file history…
Monitor This Case
Get email alerts when status or documents change.
Order Certified Copies
Most orders are placed with the USPTO same day — all within 24 business hours.
Order via The Patent Place →
Pre-filled with this patent's details
Quick Facts
Patent No.
US 10,942,828
App. No.
16/270,048
Granted
Mar 9, 2021
Kind
B2
Abstract

This application relates to distributed storage, and in particular, to a distributed shard storage technology. In a method for storing data shards in a distributed storage system, M data nodes on which to-be-stored data will be stored are determined, N replicas of the to-be-stored data are obtained, and each of the N replicas is sharded into X data shards in a same sharding mode. Then the to-be-stored data is stored on the M storage nodes, that is, N replicas of each of the X data shards are respectively stored on N storage nodes, and a quantity of data shards whose data shard replicas are stored on same N storage nodes is P or P+1, where P is an integer quotient of X divided by C M N .

Claims (41)

1. A method for storing data shards, wherein the method comprises:

determining M storage nodes on which data will be stored;

obtaining N replicas of the data, wherein the N replicas comprise original data of the data and N−1 pieces of backup data of the original data, and sharding each of the N replicas into X data shards in a same sharding mode, so that each data shard has N data shard replicas, wherein N is less than or equal to M, and wherein N and M are positive integers, wherein sharding each of the N replicas into X data shards in the same sharding mode comprises:

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N , the quantity X of the shards of the data is equal to a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode; or

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N ,

obtaining the quantity X of the shards of the data according to the optimal shard base Y, wherein the quantity X of the shards of the data is less than a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode; and

storing the N replicas of the data on the M storage nodes, wherein the N data shard replicas of each of the X data shards are respectively stored on N storage nodes in the M storage nodes, and a quantity of data shards whose data shard replicas are stored on same N storage nodes is P or P+1, wherein P is an integer quotient of X divided by C M N , wherein C M N is a total number of storage node combinations of selecting N storage nodes from the M data storage nodes.

2. The method according to claim 1 , wherein the coefficient K is determined according to a load balancing requirement of a distributed storage system, and a degree of load balancing of the data is higher as the value of K increases.

3. The method according to claim 1 , wherein obtaining the quantity N of the replicas of the data specifically comprises:

determining the quantity N of the replicas of the data according to a security requirement of the data, wherein if the quantity N of the replicas is greater, a higher security requirement of the data can be satisfied.

4. The method according to claim 1 , wherein the quantity X of the data shards is determined according to a load balancing requirement of a distributed storage system, and a degree of load balancing of the data is higher as the value of X increases.

5. The method according to claim 1 , wherein storing the N replicas of the data on the M storage nodes specifically comprises:

determining C M N data node combinations for selecting N data nodes from the M data nodes when storing the N replicas of the data on the M data nodes;

determining the quotient P obtained by dividing the quantity X of the shards by C M N , and a remainder Q; and

selecting, from the C M N data node combinations, Q data node combinations for storing P+1 data shards, and using remaining C M N −Q data node combinations for storing P data shards, wherein the N replicas of each data shard are respectively stored on N different data nodes in a data node combination.

6. A distributed storage device, wherein the device is used in a distributed storage system comprising at least two storage nodes, and configured to determine a policy for storing shards of data, and the device comprises:

a memory that stores at least one instruction; and

a processor coupled with the memory, wherein the processor executes the at least one instruction stored in the memory to perform operations, comprising:

determining M storage nodes on which data will be stored,

obtaining N replicas of the data, wherein the N replicas comprise original data of the data and N−1 pieces of backup data of the original data, and sharding each of the N replicas into X data shards in a same sharding mode, so that each data shard has N data shard replicas, wherein N is less than or equal to M, and wherein N and M are positive integers, wherein sharding each of the N replicas into X data shards in the same sharding mode comprises:

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N , the quantity X of the shards of the data is equal to a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode; or

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N ,

obtaining the quantity X of the shards of the data according to the optimal shard base Y, wherein the quantity X of the shards of the data is less than a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode, and

storing the N replicas of the data on the M storage nodes, wherein the N data shard replicas of each of the X data shards are respectively stored on N storage nodes in the M storage nodes, a quantity of data shards whose data shard replicas are stored on same N storage nodes is P or P+1, wherein P is an integer quotient of X divided by C M N , wherein C M N is a total number of storage node combinations of selecting N storage nodes from the M data storage nodes.

7. A distributed storage system, comprising:

at least two storage nodes; and

at least one management device communicatively coupled with the at least two storage nodes, wherein the management device is configured to determine a policy for storing shards of data, and wherein the management device comprises:

a memory that stores at least one instruction, and

a processor coupled with the memory, wherein the processor executes the at least one instruction stored in the memory to perform operations, comprising:

determining M storage nodes on which data will be stored,

obtaining N replicas of the data, wherein the N replicas comprise original data of the data and N−1 pieces of backup data of the original data, and sharding each of the N replicas into X data shards in a same sharding mode, so that each data shard has N data shard replicas, wherein N is less than or equal to M, and wherein N and M are positive integers, wherein sharding each of the N replicas into X data shards in the same sharding mode comprises:

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N , the quantity X of the shards of the data is equal to a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode; or

determining an optimal shard base Y of the data according to the quantity N of the replicas and the quantity M of the storage nodes, wherein the optimal shard base Y is equal to C M N ,

obtaining the quantity X of the shards of the data according to the optimal shard base Y, wherein the quantity X of the shards of the data is less than a product of the optimal shard base Y and a coefficient K, and K is an integer greater than or equal to 1, and

sharding each of the N replicas into the X data shards in the same sharding mode, and

storing the N replicas of the data on the M storage nodes, wherein the N data shard replicas of each of the X data shards are respectively stored on N storage nodes in the M storage nodes, a quantity of data shards whose data shard replicas are stored on same N storage nodes is P or P+1, wherein P is an integer quotient of X divided by C M N , wherein C M N is a total number of storage node combinations of selecting N storage nodes from the M data storage nodes.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Feb 13, 2019
From: WANG, HUAQIONG; GAO, CHAO
To: HUAWEI TECHNOLOGIES CO., LTD.
Reel/Frame 048324/0816 →
Priority Claims (1)
CN 201610659118.4 · Aug 10, 2016 · national
Continuity (2)
Continuation PCTCN2017079971 · Apr 10, 2017
Related Publication 20190171537A1 · Jun 6, 2019