IP Library › Granted Patent US 10,496,614
Granted Patent B2
US 10,496,614 · App. 15/268,318 · Granted Dec 3, 2019

DDL processing in shared databases

Inventors: Wei-Ming Hu (Palo Alto, CA); Mark Dilman (Sunnyvale, CA); Leonid Novak (Castro Valley, CA); Stephen Ball (Dublin, CA); Ghazi Nourdine Benadjaoud (Pleasanton, CA)
Assignee: Oracle International Corporation
G06F16/213G06F16/217G06F16/221G06F16/2272G06F16/2282G06F16/248G06F16/2455G06F16/2471G06F16/252G06F16/27G06F16/278
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,496,614
App. No.
15/268,318
Granted
Dec 3, 2019
Kind
B2
Abstract

Techniques are provided for creating, organizing, and maintaining a sharded database. A sharded database can be created using user-defined sharding, system-managed sharding, or composite sharding. The sharded database is implemented with relational database techniques. The techniques described can be used for load distribution, organization, query processing, and schema propagation in a sharded database.

Claims (82)

1. A method, comprising:

maintaining, by a shard catalogue, schema data that indicates a schema that is used by each shard, of a plurality of shards, of a sharded database;

wherein each shard, of the plurality of shards, is a distinct relational database server;

wherein the shard catalogue is a special database that is used to store configuration data for the sharded database;

receiving schema modification instructions for modifying the schema;

in response to receiving the schema modification instructions, automatically performing:

updating, at the shard catalogue, the schema data based on the schema modification instructions; and

causing the respective relational database servers of all shards of the plurality of shards to modify the schema by, for each particular shard of the plurality of shards:

creating a database connection with the respective relational database server of the particular shard;

sending the schema modification instructions to the respective relational database server of the particular shard; and

executing the schema modification instructions within the respective relational database server of the particular shard.

2. The method of claim 1 , wherein each relational database server of the sharded database does not share processors, memory, or disk storage with other relational database servers of the sharded database.

3. The method of claim 1 , wherein the schema modification instructions are written in a data definition language (DDL).

4. The method of claim 1 , further comprising sending the schema modification instructions from the shard catalogue to a shard director.

5. The method of claim 4 , further comprising:

for each particular shard of the plurality of shards:

after executing the schema modification instructions on the particular shard, sending, from the particular shard to the shard director, a notification that indicates a status of executing the schema modification instructions on the particular shard; and

updating, by the shard director, the schema data to indicate the status of executing the schema modification instructions on the particular shard.

6. A method, comprising:

maintaining, by a shard catalogue, schema data that indicates a schema that is used by each shard, of a plurality of shards, of a sharded database;

wherein each shard, of the plurality of shards, is a distinct relational database server;

wherein the shard catalogue is a special database that is used to store configuration data for the sharded database;

receiving schema modification instructions for modifying the schema;

in response to receiving the schema modification instructions, automatically performing:

updating, at the shard catalogue, the schema data based on the schema modification instructions;

sending the schema modification instructions from the shard catalogue to a shard director, for each particular shard of the plurality of shards:

after executing the schema modification instructions on the particular shard, sending, from the particular shard to the shard director, a notification that indicates a status of executing the schema modification instructions on the particular shard; and

updating, by the shard director, the schema data to indicate the status of executing the schema modification instructions on the particular shard;

in response to detecting, by the shard director, that a first shard is unavailable, deferring the execution of the schema modification instructions on the first shard; and

causing the respective relational database servers of all shards of the plurality of shards to modify the schema by, for each particular shard of the plurality of shards:

creating a database connection with the respective relational database server of the particular shard;

sending the schema modification instructions to the respective relational database server of the particular shard;

executing the schema modification instructions within the respective relational database server of the particular shard.

7. The method of claim 6 , further comprising:

in response to detecting, by the shard director, that the first shard is available, performing the execution of the schema modification instructions on the first shard.

8. The method of claim 1 , wherein the schema modification instructions comprise instructions for creating a sharded table; and wherein executing the schema modification instructions on the particular shard comprises creating, by the particular shard, a partition of the sharded table to be stored in a chunk of data on the particular shard without creating a complete sharded table on the particular shard.

9. The method of claim 1 , further comprising:

receiving, by a shard director, notification statuses from the plurality of shards regarding a workload of each particular shard; and

for each particular shard of the plurality of shards:

performing the steps of creating a database connection, sending the schema modification instructions to the particular shard, and executing the schema modification instructions only when the workload of each particular shard is below a threshold workload setting.

10. The method of claim 1 , wherein the causing all shards of the plurality of shards to modify the schema is performed in parallel for each shard of the plurality of shards.

11. One or more non-transitory computer-readable media storing instructions, wherein the instructions include:

instructions which, when executed by one or more hardware processors, cause maintaining, by a shard catalogue, schema data that indicates a schema that is used by each shard, of a plurality of shards, of a sharded database;

wherein each shard, of the plurality of shards, is a distinct relational database server;

wherein the shard catalogue is a special database that is used to store configuration data for the sharded database;

instructions which, when executed by one or more hardware processors, cause receiving schema modification instructions for modifying the schema;

instructions which, when executed by one or more hardware processors, cause in response to receiving the schema modification instructions, automatically performing:

instructions which, when executed by one or more hardware processors, cause updating, at the shard catalogue, the schema data based on the schema modification instructions; and

instructions which, when executed by one or more hardware processors, cause causing the respective relational database servers of all shards of the plurality of shards to modify the schema by, for each particular shard of the plurality of shards:

instructions which, when executed by one or more hardware processors, cause creating a database connection with the respective relational database server of the particular shard;

instructions which, when executed by one or more hardware processors, cause sending the schema modification instructions to the respective relational database server of the particular shard; and

instructions which, when executed by one or more hardware processors, cause executing the schema modification instructions within the respective relational database server of the particular shard.

12. The one or more non-transitory computer-readable media of claim 11 , wherein each relational database server of the sharded database does not share processors, memory, or disk storage with other relational database servers of the sharded database.

13. The one or more non-transitory computer-readable media of claim 11 , wherein the schema modification instructions are written in a data definition language (DDL).

14. The one or more non-transitory computer-readable media of claim 11 , further comprising instructions which, when executed by one or more hardware processors, cause sending the schema modification instructions from the shard catalogue to a shard director.

15. The one or more non-transitory computer-readable media of claim 14 , further comprising:

for each particular shard of the plurality of shards:

instructions which, when executed by one or more hardware processors, cause, after executing the schema modification instructions on the particular shard, sending, from the particular shard to the shard director, a notification that indicates a status of executing the schema modification instructions on the particular shard; and

instructions which, when executed by one or more hardware processors, cause, updating, by the shard director, the schema data to indicate the status of executing the schema modification instructions on the particular shard.

16. The one or more non-transitory computer-readable media storing instructions, wherein the instructions include:

instructions which, when executed by one or more hardware processors, cause maintaining, by a shard catalogue, schema data that indicates a schema that is used by each shard, of a plurality of shards, of a sharded database;

wherein each shard, of the plurality of shards, is a distinct relational database server;

wherein the shard catalogue is a special database that is used to store configuration data for the sharded database;

instructions which, when executed by one or more hardware processors, cause receiving schema modification instructions for modifying the schema;

instructions which, when executed by one or more hardware processors, cause in response to receiving the schema modification instructions, automatically performing:

instructions which, when executed by one or more hardware processors, cause updating, at the shard catalogue, the schema data based on the schema modification instructions;

instructions which, when executed by one or more hardware processors, cause sending the schema modification instructions from the shard catalogue to a shard director, for each particular shard of the plurality of shards:

instructions which, when executed by one or more hardware processors, cause, after executing the schema modification instructions on the particular shard, sending, from the particular shard to the shard director, a notification that indicates a status of executing the schema modification instructions on the particular shard; and

instructions which, when executed by one or more hardware processors, cause, updating, by the shard director, the schema data to indicate the status of executing the schema modification instructions on the particular shard;

instructions which, when executed by one or more hardware processors, cause, in response to detecting, by the shard director, that a first shard is unavailable, deferring the execution of the schema modification instructions on the first shard; and

instructions which, when executed by one or more hardware processors, cause causing the respective relational database servers of all shards of the plurality of shards to modify the schema by, for each particular shard of the plurality of shards:

instructions which, when executed by one or more hardware processors, cause creating a database connection with the respective relational database server of the particular shard;

instructions which, when executed by one or more hardware processors, cause sending the schema modification instructions to the respective relational database server of the particular shard; and

instructions which, when executed by one or more hardware processors, cause executing the schema modification instructions within the respective relational database server of the particular shard.

17. The one or more non-transitory computer-readable media of claim 16 , further comprising:

instructions which, when executed by one or more hardware processors, cause, in response to detecting, by the shard director, that the first shard is available, performing the execution of the schema modification instructions on the first shard.

18. The one or more non-transitory computer-readable media of claim 11 , wherein the schema modification instructions comprise instructions for creating a sharded table; and wherein the instructions for executing the schema modification instructions on the particular shard comprises instructions which, when executed by one or more hardware processors, cause, creating, by the particular shard, a partition of the sharded table to be stored in a chunk of data on the particular shard without creating a complete sharded table on the particular shard.

19. The one or more non-transitory computer-readable media of claim 11 , further comprising:

instructions which, when executed by one or more hardware processors, cause receiving, by a shard director, notification statuses from the plurality of shards regarding a workload of each particular shard; and

for each particular shard of the plurality of shards:

instructions which, when executed by one or more hardware processors, cause performing the steps of creating a database connection, sending the schema modification instructions to the particular shard, and executing the schema modification instructions only when the workload of each particular shard is below a threshold workload setting.

20. The one or more non-transitory computer-readable media of claim 11 wherein the instructions which, when executed by one or more hardware processors, cause causing all shards of the plurality of shards to modify the schema is performed in parallel for each shard of the plurality of shards.

Assignments (1)
ASSIGNMENT OF ASSIGNOR'S INTEREST Recorded Sep 16, 2016
From: HU, WEI-MING; DILMAN, MARK; NOVAK, LEONID; BALL, STEPHEN; BENADJAOUD, GHAZI NOURDINE
To: ORACLE INTERNATIONAL CORPORATION
Reel/Frame 039770/0383 →
Continuity (2)
Provisional Application 62238193 · Oct 7, 2015
Related Publication 20170103092A1 · Apr 13, 2017
Cited By (1)
US 12,608,355