O particionamento e o bucketing no ApsaraDB for SelectDB dividem os dados em intervalos gerenciáveis e os distribuem entre os nós para otimizar o armazenamento e a computação.
Visão geral
Para armazenar e processar grandes volumes de dados com eficiência, o ApsaraDB for SelectDB divide os dados em partições e os distribui pelo cluster para processamento paralelo.
Todos os modelos de dados no ApsaraDB for SelectDB suportam os dois níveis de particionamento de dados a seguir:
-
Um nível: os dados são particionados em apenas um nível.
Ao criar uma tabela sem especificar partições, o ApsaraDB for SelectDB cria uma partição padrão transparente para você. Nesse caso, há suporte apenas para bucketing.
-
Dois níveis: os dados são particionados em dois níveis.
O primeiro nível é a partição, que suporta particionamento por intervalo (range) e por lista (list).
O segundo nível é o bucket, também conhecido como tablet, que suporta particionamento por hash.
Particionamento
As partições dividem os dados em intervalos distintos, semelhante à divisão de uma tabela em várias sub-tabelas para facilitar o gerenciamento. Ao utilizar partições, observe os seguintes pontos:
Especifique uma ou mais colunas como colunas de chave de partição. As colunas de chave de partição devem ser colunas-chave.
Os valores da chave de partição devem estar entre aspas duplas ("), independentemente do tipo de dado da coluna.
Teoricamente, não há limite para o número de partições criadas.
Se você criar uma tabela sem definir partições, o sistema gera automaticamente uma partição com o mesmo nome da tabela, contendo todos os seus dados. Essa partição é invisível e não pode ser excluída ou modificada.
Ao criar uma partição, o intervalo definido não pode sobrepor o intervalo de outra partição existente.
Particionamento por intervalo (Range)
Colunas de tempo são frequentemente usadas como chaves de partição por intervalo para facilitar o gerenciamento de dados novos e históricos. Nas partições por intervalo, especifique apenas o limite superior executando a instrução VALUES LESS THAN (...). O sistema utiliza o limite superior da partição anterior como limite inferior da partição atual, criando um intervalo fechado à esquerda e aberto à direita. Também é possível definir ambos os limites (superior e inferior) executando a instrução VALUES [...], criando igualmente um intervalo fechado à esquerda e aberto à direita.
Particionamento por coluna única
O exemplo a seguir demonstra como os intervalos de partição se comportam ao usar a instrução VALUES LESS THAN (...) para criar ou excluir partições:
-
Crie uma tabela chamada test_table.
CREATE TABLE IF NOT EXISTS test_db.test_table ( `user_id` LARGEINT NOT NULL COMMENT "The user ID", `date` DATE NOT NULL COMMENT "The date on which data is imported to the table", `timestamp` DATETIME NOT NULL COMMENT "The time when data is imported to the table", `city` VARCHAR(20) COMMENT "The city in which the user resides", `age` SMALLINT COMMENT "The age of the user", `sex` TINYINT COMMENT "The gender of the user", `last_visit_date` DATETIME REPLACE DEFAULT "1970-01-01 00:00:00" COMMENT "The last time when the user paid a visit", `cost` BIGINT SUM DEFAULT "0" COMMENT "The amount of money that the user spends", `max_dwell_time` INT MAX DEFAULT "0" COMMENT "The maximum dwell time of the user", `min_dwell_time` INT MIN DEFAULT "99999" COMMENT "The minimum dwell time of the user" )ENGINE=OLAP AGGREGATE KEY(`user_id`, `date`, `timestamp`, `city`, `age`, `sex`) PARTITION BY RANGE(`date`) ( PARTITION `p201701` VALUES LESS THAN ("2017-02-01"), PARTITION `p201702` VALUES LESS THAN ("2017-03-01"), PARTITION `p201703` VALUES LESS THAN ("2017-04-01") ) DISTRIBUTED BY HASH(`user_id`) BUCKETS 16;Após a criação da tabela
test_table, as três partições a seguir são geradas automaticamente:p201701: [MIN_VALUE, 2017-02-01) p201702: [2017-02-01, 2017-03-01) p201703: [2017-03-01, 2017-04-01) -
Execute a instrução
ALTER TABLE test_db.test_table ADD PARTITION p201705 VALUES LESS THAN ("2017-06-01");para criar uma partição chamadap201705. O código de exemplo a seguir mostra o resultado do particionamento:p201701: [MIN_VALUE, 2017-02-01) p201702: [2017-02-01, 2017-03-01) p201703: [2017-03-01, 2017-04-01) p201705: [2017-04-01, 2017-06-01) -
Execute a instrução
ALTER TABLE test_db.test_table DROP PARTITION p201703;para excluir a partiçãop201703. O código de exemplo a seguir mostra o resultado do particionamento:p201701: [MIN_VALUE, 2017-02-01) p201702: [2017-02-01, 2017-03-01) p201705: [2017-04-01, 2017-06-01)ImportanteNo exemplo anterior, após a exclusão da partição p201703, os intervalos das partições p201702 e p201705 permanecem inalterados. No entanto, o intervalo [2017-03-01,2017-04-01) entre eles fica vago. Os dados existentes nesse intervalo também são excluídos. Nesse cenário, se os dados a serem importados estiverem dentro do intervalo vago, a importação falhará.
-
Exclua a partição p201702. O código de exemplo a seguir mostra o resultado do particionamento:
p201701: [MIN_VALUE, 2017-02-01) p201705: [2017-04-01, 2017-06-01)O intervalo vago passa a ser [2017-02-01,2017-04-01).
-
Execute a instrução
p201702newpara criar uma partição. O código de exemplo a seguir mostra o resultado do particionamento:p201701: [MIN_VALUE, 2017-02-01) p201702new: [2017-02-01, 2017-03-01) p201705: [2017-04-01, 2017-06-01)O intervalo vago passa a ser [2017-03-01,2017-04-01).
-
Exclua a partição p201701 e execute a instrução
p201612para criar uma partição. O código de exemplo a seguir mostra o resultado do particionamento:p201612: [MIN_VALUE, 2017-01-01) p201702new: [2017-02-01, 2017-03-01) p201705: [2017-04-01, 2017-06-01)Os intervalos vagos passam a ser [2017-01-01,2017-02-01) e [2017-03-01,2017-04-01).
Conforme demonstrado no exemplo anterior, os intervalos das partições existentes permanecem inalterados após a exclusão de partições, mas podem surgir intervalos vagos. Ao usar a instrução VALUES LESS THAN (...) para criar uma nova partição, o limite inferior deve ser adjacente ao limite superior da partição anterior.
Particionamento por múltiplas colunas
É possível particionar dados com base em várias colunas. Exemplo:
PARTITION BY RANGE(`date`, `id`)
(
PARTITION `p201701_1000` VALUES LESS THAN ("2017-02-01", "1000"),
PARTITION `p201702_2000` VALUES LESS THAN ("2017-03-01", "2000"),
PARTITION `p201703_all` VALUES LESS THAN ("2017-04-01")
)
Neste exemplo, as colunas date e id são definidas como colunas de chave de partição. A coluna date é do tipo DATE e a coluna id é do tipo INT. O código de exemplo a seguir mostra o resultado do particionamento:
* p201701_1000: [(MIN_VALUE, MIN_VALUE), ("2017-02-01", "1000") )
* p201702_2000: [("2017-02-01", "1000"), ("2017-03-01", "2000") )
* p201703_all: [("2017-03-01", "2000"), ("2017-04-01", MIN_VALUE))
Na última partição, apenas o valor da coluna date foi especificado. Por padrão, MIN_VALUE é usado como valor da coluna id. Durante a inserção de dados, o sistema compara os dados com os valores da chave de partição especificados sequencialmente para determinar a partição de destino. O código de exemplo a seguir ilustra esse comportamento:
* Data --> Partition
* 2017-01-01, 200 --> p201701_1000
* 2017-01-01, 2000 --> p201701_1000
* 2017-02-01, 100 --> p201701_1000
* 2017-02-01, 2000 --> p201702_2000
* 2017-02-15, 5000 --> p201702_2000
* 2017-03-01, 2000 --> p201703_all
* 2017-03-10, 1 --> p201703_all
* 2017-04-01, 1000 --> Failed to be imported.
* 2017-05-01, 1000 --> Failed to be imported.
Particionamento por lista (List)
O particionamento por lista suporta os seguintes tipos de dados para colunas de chave de partição: BOOLEAN, TINYINT, SMALLINT, INT, BIGINT, LARGEINT, DATE, DATETIME, CHAR e VARCHAR. Os dados são importados para uma partição somente quando correspondem a um dos valores de enumeração dessa partição.
Especifique os valores de enumeração contidos em cada partição executando a instrução VALUES IN (...).
Particionamento por coluna única
O exemplo a seguir demonstra como as partições se comportam ao usar a instrução VALUES IN (...) para criar ou excluir partições:
-
Crie uma tabela chamada test_table1.
CREATE TABLE IF NOT EXISTS test_db.example_list_tbl1 ( `user_id` LARGEINT NOT NULL COMMENT "The user ID", `date` DATE NOT NULL COMMENT "The date on which data is imported to the table", `timestamp` DATETIME NOT NULL COMMENT "The time when data is imported to the table", `city` VARCHAR(20) NOT NULL COMMENT "The city in which the user resides", `age` SMALLINT COMMENT "The age of the user", `sex` TINYINT COMMENT "The gender of the user", `last_visit_date` DATETIME REPLACE DEFAULT "1970-01-01 00:00:00" COMMENT "The last time when the user paid a visit", `cost` BIGINT SUM DEFAULT "0" COMMENT "The amount of money that the user spends", `max_dwell_time` INT MAX DEFAULT "0" COMMENT "The maximum dwell time of the user", `min_dwell_time` INT MIN DEFAULT "99999" COMMENT "The minimum dwell time of the user" ) ENGINE=olap AGGREGATE KEY(`user_id`, `date`, `timestamp`, `city`, `age`, `sex`) PARTITION BY LIST(`city`) ( PARTITION `p_cn` VALUES IN ("Beijing", "Shanghai", "Hong Kong"), PARTITION `p_usa` VALUES IN ("New York", "San Francisco"), PARTITION `p_jp` VALUES IN ("Tokyo") ) DISTRIBUTED BY HASH(`user_id`) BUCKETS 16;Após a criação da tabela test_table1, as três partições a seguir são geradas automaticamente:
p_cn: ("Beijing", "Shanghai", "Hong Kong") p_usa: ("New York", "San Francisco") p_jp: ("Tokyo") -
Execute a instrução
p_ukpara criar uma partição. O código de exemplo a seguir mostra o resultado do particionamento:p_cn: ("Beijing", "Shanghai", "Hong Kong") p_usa: ("New York", "San Francisco") p_jp: ("Tokyo") p_uk: ("London") -
Exclua a partição p_jp. O código de exemplo a seguir mostra o resultado do particionamento:
p_cn: ("Beijing", "Shanghai", "Hong Kong") p_usa: ("New York", "San Francisco") p_uk: ("London")
Particionamento por múltiplas colunas
É possível particionar dados com base em várias colunas. Exemplo:
PARTITION BY LIST(`id`, `city`)
(
PARTITION `p1_city` VALUES IN (("1", "Beijing"), ("1", "Shanghai")),
PARTITION `p2_city` VALUES IN (("2", "Beijing"), ("2", "Shanghai")),
PARTITION `p3_city` VALUES IN (("3", "Beijing"), ("3", "Shanghai"))
)
Neste exemplo, as colunas id e city são definidas como colunas de chave de partição. A coluna id é do tipo INT e a coluna city é do tipo VARCHAR. O código de exemplo a seguir mostra o resultado do particionamento:
* p1_city: [("1", "Beijing"), ("1", "Shanghai")]
* p2_city: [("2", "Beijing"), ("2", "Shanghai")]
* p3_city: [("3", "Beijing"), ("3", "Shanghai")]
Durante a inserção de dados, o sistema compara os dados com os valores da chave de partição especificados sequencialmente para determinar a partição de destino. O código de exemplo a seguir ilustra esse comportamento:
* Data ---> Partition
* 1, Beijing ---> p1_city
* 1, Shanghai ---> p1_city
* 2, Shanghai ---> p2_city
* 3, Beijing ---> p3_city
* 1, Tianjin ---> Failed to be imported.
* 4, Beijing ---> Failed to be imported.
Bucketing
Os dados são distribuídos entre os buckets com base nos valores de hash das colunas de bucket especificadas.
Se houver partições, a instrução DISTRIBUTED... define as regras de divisão dos dados dentro de cada partição. Caso contrário, a instrução define as regras para dividir todos os dados da tabela.
Especifique várias colunas como colunas de bucket. Nos modelos Aggregate ou Unique, as colunas de bucket devem ser colunas-chave. No modelo Duplicate, elas podem ser colunas-chave ou colunas de valor. As colunas de bucket podem ser iguais ou diferentes das colunas de chave de partição.
-
Ao escolher as colunas de bucket, equilibre o throughput de consulta e a concorrência de consulta.
Especificar múltiplas colunas de bucket resulta em uma distribuição de dados mais uniforme. Contudo, se uma consulta não incluir condições equivalentes para todas as colunas de bucket, o sistema varrerá todos os buckets. Isso aumenta o throughput da consulta e reduz a latência, sendo adequado para cenários de alto throughput e baixa concorrência.
Especificar apenas uma ou poucas colunas de bucket permite que o sistema varra apenas um bucket para consultas pontuais. Quando múltiplas consultas pontuais são executadas simultaneamente, elas podem acessar buckets diferentes sem interferência nas operações de I/O, especialmente se os buckets estiverem em discos distintos. Esse cenário é ideal para consultas pontuais de alta concorrência.
Teoricamente, não há limite para o número de buckets criados.
Melhores práticas
Recomendações para configurar partições e buckets**
O número total de buckets em uma tabela segue a fórmula: Total de buckets = Número de partições × Número de buckets por partição.
Mantendo-se as configurações do cluster inalteradas, recomenda-se que o número de buckets em uma partição seja ligeiramente superior ao total de discos no cluster.
Armazene de 1 a 10 GB de dados por bucket. Poucos dados em um bucket reduzem a eficácia da agregação e aumentam a sobrecarga de gerenciamento de metadados. Muitos dados tornam a migração e a reposição de réplicas mais lentas, além de elevar o custo de novas tentativas em operações no nível do bucket, como alterações de esquema ou rollups.
Caso não seja possível equilibrar o tamanho dos dados por bucket e a quantidade de buckets, priorize o tamanho dos dados por bucket.
Na criação da tabela, define-se o mesmo número de buckets para cada partição. Porém, ao executar a instrução ADD PARTITION para criar uma partição dinamicamente, é possível especificar separadamente o número de buckets na nova partição. Use esse recurso para lidar com redução ou expansão de dados.
Não é possível alterar o número de buckets de uma partição após sua criação. Planeje a contagem de buckets considerando futuras expansões do cluster. Por exemplo, se o cluster tiver três máquinas com um disco cada e a contagem de buckets for definida como 3 ou menos, adicionar mais máquinas não melhorará a concorrência.
A tabela a seguir apresenta recomendações de partições e buckets para um cluster com 10 backends, cada um com um disco.
|
Tamanho da tabela |
500 MB |
5 GB |
50 GB |
500 GB |
5 TB |
|
Partições |
Nenhuma partição necessária. |
Nenhuma partição necessária. |
Nenhuma partição necessária. |
Cada partição tem 50 GB. |
Cada partição tem 50 GB. |
|
Buckets |
A tabela contém de quatro a oito buckets. |
A tabela contém de 8 a 16 buckets. |
A tabela contém 32 buckets. |
Cada partição contém de 16 a 32 buckets. |
Cada partição contém de 16 a 32 buckets. |
Execute a instrução SHOW DATA; para consultar o tamanho de uma tabela.
Configuração e uso do método de distribuição aleatória
Para dados detalhados que não exigem agregação ou atualizações, crie uma tabela usando o modelo Duplicate com o método de distribuição aleatória. Exemplo:
CREATE TABLE IF NOT EXISTS test.example_tbl
(
`timestamp` DATETIME NOT NULL COMMENT "The time when the log was generated",
`type` INT NOT NULL COMMENT "The type of the log",
`error_code` INT COMMENT "The error code",
`error_msg` VARCHAR(1024) COMMENT "The error message",
`op_id` BIGINT COMMENT "The owner ID",
`op_time` DATETIME COMMENT "The time when the error was handled"
)
DUPLICATE KEY(`timestamp`, `type`, `error_code`)
DISTRIBUTED BY RANDOM BUCKETS 16;
Tabelas que usam o modelo de chave Duplicate não contêm colunas com tipo de agregação REPLACE. Defina o modo de bucketing de dados da tabela como RANDOM para evitar desequilíbrios graves de dados. Durante a importação, um único job de importação grava dados em um bucket aleatório da partição.
Com o bucketing RANDOM, não é possível consultar buckets específicos por valores de coluna de bucket, pois nenhuma coluna de bucket é especificada. O sistema varre todos os buckets na partição correspondente. Essa abordagem é adequada para consultas e análises agregadas de tabela completa, e não para consultas pontuais de alta concorrência.
Para tabelas do modelo Duplicate que usam distribuição aleatória, ative o modo de importação de bucket único definindo o parâmetro load_to_single_tablet como true (padrão: false). Nesse modo, cada job de importação grava dados em apenas um bucket por partição, o que melhora a concorrência e o throughput de importação, reduz a amplificação de escrita por compactação e ajuda a manter a estabilidade do cluster.
Cenários de uso combinado de partições e buckets
Se uma tabela contiver colunas de dimensão temporal ou colunas de dimensão com valores ordenados, essas colunas podem servir como chaves de partição. Avalie a granularidade do particionamento com base na frequência de importação e no volume de dados a ser armazenado em cada partição.
Para excluir dados históricos e reter apenas os dados dos últimos N dias, utilize o particionamento composto para remover partições antigas. Alternativamente, execute a instrução DELETE para apagar os dados de uma partição específica.
Para evitar desequilíbrio de dados, especifique separadamente o número de buckets para cada partição. Por exemplo, em cenários com particionamento diário onde o volume de dados varia significativamente, personalize a contagem de buckets para cada partição. Especifique colunas de bucket fáceis de identificar e que permitam uma distribuição uniforme dos dados.