Differences

This shows you the differences between two versions of the page.

Link to this comparison view

Both sides previous revision Previous revision
design:partitioning:intro [2009/07/23 00:45]
admin
— (current)
Line 1: Line 1:
-====== Partitioning ====== 
  
-==== Partitioning ======= 
- 
-Partitioning is a common RDBMS strategy that turns a large set of data into smaller, logical subsets. The most common objectives for partitioning databases are to make application queries perform better and to make the data easier to administrate. 
- 
-Omnidex supports partitioning in three ways: 
- 
-physical partitioning of data into two or more separate partitions ​ 
-partitioning of indexes using composite MDK indexes ​ 
-partitioning of indexes using composite ASK indexes ​ 
-All three partitioning concepts can be used together or separately and are discussed in detail in this section. 
- 
-What if my data is already partitioned? ​ 
-  
- 
-Advantages 
-Partitioned data can return extraordinary performance on extremely large databases. A partitioned table is broken up into smaller pieces, making the individual table sizes smaller. For example, a table contaning 10 billion rows partitioned into 10 partitions, creates 10 one billion row tables. Searching one billion rows is considerably faster than searching 10 billion rows, and even faster when installed with Omnidex indexes. 
- 
-Physical index files are smaller, narrowing the number of index key values Omnidex must look through to qualify data. 
- 
-Index build time can be shorter. For example, a partition can be made up of updates for a given period of time. The new partition can be appended to the end of the existing partitions. The indexes for this new partition can be built without the need to rebuild the indexes on the other partitions. 
- 
-Disadvantages 
-Advance analysis is required to determine the best partitioning strategy, and indeed if the data should be partitioned at all. 
- 
-Continual analysis may be required to determine if the current partitioning strategy continues to be the best strategy. 
- 
-Since the data is physically separated into different tables, updates to the data must be made to the respective partition. At this time, Omnidex does not handle this automatically. 
- 
-  
- 
-Limitations 
-Only child tables can be partitioned. ​ 
-Only one partitioned child table per query is allowed. ​ 
-Select items containing multiple aggregations in a single expression are not supported. ​ 
-Left outer joins TO a partitioned table are not allowed. However, left outer joins FROM a partitioned table are allowed. ​ 
-Updates to partitioned data must be made to the individual partition. Omnidex does not handle this automatically. See Updating Partitioned Tables (below) for more information. ​ 
-MDK Composite keys are limited to 240 bytes. ​ 
-Aggregation indexes are limited to 32 values in an IN clause. ​ 
-Since each partition will have its own set of indexes and the number of index files is limited to 255 physical files, the number of partitions can be limited. ​ 
-  
- 
-Updating Partitioned Tables 
-Applications that update partitioned tables must take some extra steps to make sure the updates are made to the correct partition. Unlike select statements, insert, update and delete statements must reference the individual partition, not the unioned table. 
- 
-For example, if the prospects table is partitioned into 5 partitions, where the first partition contains only prospects from the state of California, the second partition contains only prospects from the states of New York and New Jersery, and so one, the application would have to determine which partition to perform the update against according to the state. 
- 
-If a prospect from the state of California is inserted into the second partition, which contains only prospects from New York and New Jersey, the new record would never be qualified by any select statement. 
- 
-Therefore, the update application must determine which partition should be updated. 
- 
-  
- 
-What if my data is already partitioned?​ 
-If the data is already partitioned,​ analysis is still the first and most important step. 
- 
-Is there a partition qualifier? ​ 
-Does the partition structure meet the needs of all of my queries? ​ 
-Are the partitions relatively close in size (number of rows)? ​ 
-The Analysis topic will help to answer the questions and help you determine how to proceed. 
- 
-Multiple-predicate partitioning 
- 
-- Union View and Grid partitioning can now use multiple predicates 
-- Predicates must be AND’d 
-- Addin will determine which partitions to use based on the extent that one or more predicates is met. 
- 
- 
- 
- 
- 
- 
- 
- 
-Example of multiple-predicate partitioning 
- 
- 
-node GRID01 
- ​database LIST filedsn "​list01.dsn"​ local cache 
-  cluster 
-   table LIST partition by "STATE in ('​MA','​ME','​NH','​NY','​PR','​RI','​VT'​) and  
-                            ZIP between '​00000'​ and '​05999'"​ 
- 
-node GRID02 
- ​database LIST filedsn "​list02.dsn"​ local cache 
-  cluster 
-   table LIST partition by "STATE in ('​CT','​NJ','​NY'​) and  
-                            ZIP between '​06000'​ and '​10999'"​ 
- 
-node GRID03 
- ​database LIST filedsn "​list03.dsn"​ local cache 
-  cluster 
-   table LIST partition by "STATE in ('​NY'​) and  
-                            ZIP between '​11000'​ and '​14999'"​ 
- 
- 
-  
-- Queries with criteria of STATE = ‘VT’ will hit first partition only  
-- Queries with criteria of ZIP = ‘10000’ will hit second partition only 
-- Queries with criteria of STATE = ‘NY’ will hit all three partitions. 
-- Queries with criteria of STATE = ‘NY’ and ZIP = ‘10000’ will hit the second partition only. 
- 
-  
-Database partitioning consists of dividing a large database into smaller chunks. ​ 
- 
-This is generally done for one of three reasons.  ​ 
- 
-  *  Most often, this is done to improve performance by allowing parallel processing of the partitions. ​ A query which must process the whole table can be divided into several queries that run in parallel, each evaluating their own partition.  ​ 
- 
-  * A database can also be partitioned to accommodate a large number of new entries. ​ New entries can be placed in a separate partition and indexed independently,​ avoiding the cost of re-indexing the entire database.  ​ 
- 
-  * A database can also be partitioned to segregate data, such as dividing up a database by client or region. ​ This can support security and backup strategies, and can allow different resources to be allocated to individual partitions. 
- 
-Typically, database partitioning consists of identifying one of more large tables that should be split into multiple chunks. ​ Sometimes only one table needs to be partitioned,​ but more often, a parent table is connected to multiple child tables, and the family of tables needs to be partitioned. ​ For example, a table of HOUSEHOLDS may be a parent to a table of INDIVIDUALS. ​ Furthermore,​ there may be other child tables that describe HOUSEHOLDS, and there may also be other grandchild tables that describe INDIVIDUALS. ​ As long as these are divided up in a way that preserves these parent-child relationships,​ these can be partitioned as a cluster. ​ It is important that all of the INDIVIDUALS for a given HOUSEHOLD are all found in the same partition, and the same must be true for all of the other children as well. 
- 
-With Omnidex, partitioning is typically done using a partition qualifier. ​ A partition qualifier is simply the criterion which is used to divide the database into separate partitions. ​ In the example above, the HOUSEHOLD_ID would be a good partition qualifier since it insures that all of the rows sharing a specific HOUSEHOLD_ID would be placed in the same partition. ​ Partition qualifiers can also be used to improve performance. ​ If the partition qualifier is found in the criteria of a query, it is possible to limit the number of partitions that are evaluated for that query. ​ For example, if a database is partitioned by STATE, a query that has criteria of STATE = ‘CO’ means that only those partitions with the state of Colorado must be evaluated. ​ This can provide substantial gains in performance. ​ Importantly,​ partitioning by STATE probably implies that it is also partitioned by the HOUSEHOLD_ID since households don’t usually span multiple states. ​ This allows us to insure that all the rows for a particular household are still found on the same partition while also allowing some queries to be restricted to a subset of partitions.  ​ 
-This same principle can be extended further. ​ Partitioning by the STATE code probably implies partitioning by other columns as well, such as postal, SCF or DMA codes. ​ The values for these columns are not shared across states. ​ Partitioning by state may not imply partitioning by columns like county codes or cities since those can repeat across multiple states. ​ It is worthwhile to consider which columns can be used as partition qualifiers. ​ When more columns can be identified as partition qualifiers, there is more likelihood that a given query can be restricted to a subset of partitions. 
- 
-The partition qualifier is integral to determining the number of partitions. ​ Deciding the number of partitions requires an understanding of the queries. ​ When most queries can be restricted to a few partitions, it is beneficial to have a larger number of partitions since this produces smaller partitions. ​ When a substantial percentage of queries have to scan many partitions, it is preferable to have fewer partitions. ​ There is overhead to merging the results of multiple partitions, especially when processing grouped aggregations or distinct operations across multiple partitions. ​ Generally, it is appropriate to have between eight and twenty partitions, thought this can be easily overruled by an understanding of the queries. 
- 
-====== Quick Links ====== ​ 
-{{page>:​quicklinks&​nofooter&​noeditbtn}} 
 
Back to top
design/partitioning/intro.1248309911.txt.gz · Last modified: 2012/10/26 14:25 (external edit)