Saiba como usar o conector do StarRocks.
Contexto
O StarRocks é um data warehouse de processamento massivamente paralelo (MPP) de última geração que oferece desempenho extremamente rápido em todos os cenários e uma experiência analítica unificada. O StarRocks apresenta as seguintes vantagens:
Compatibilidade com o protocolo MySQL, permitindo o uso de clientes MySQL e ferramentas comuns de Business Intelligence (BI) para conexão e análise de dados.
-
Arquitetura distribuída:
Particiona horizontalmente as tabelas de dados e as armazena com múltiplas réplicas.
Permite dimensionamento flexível do cluster e análise de até 10 petabytes (PB) de dados.
Usa framework MPP para acelerar computações paralelas.
Suporta múltiplas réplicas para garantir tolerância a falhas.
O conector do Flink armazena dados em cache e usa o Stream Load para gravá-los em lotes nas tabelas de destino. A leitura das tabelas de origem ocorre mediante busca de dados em lotes. A tabela a seguir lista as capacidades do conector do StarRocks.
|
Categoria |
Descrição |
|
Tipos suportados |
Tabelas de source, tabelas de dimensão, tabelas de destino e alvos de ingestão de dados |
|
Modo de execução |
Modo streaming e modo batch |
|
Formato de dados |
CSV |
|
Métricas específicas do conector |
Nenhuma |
|
Tipos de API |
DataStream, SQL e yaml para ingestão de dados |
|
Suporte a atualizações/exclusões em tabelas de destino |
Sim |
Pré-requisitos
Você deve ter um cluster StarRocks implantado no EMR ou um cluster autogerenciado no ECS.
Limitações
Somente o Ververica Runtime (VVR) 11.1 ou posterior oferece suporte a joins com tabelas de dimensão.
Para evitar restrições de acesso à rede, adicione as seguintes portas do cluster StarRocks a um grupo de segurança ou à lista de permissões do firewall: 9030, 8030, 8040, 9060, 8060, 9020.
SQL
Recursos
O StarRocks no E-MapReduce suporta instruções CREATE TABLE AS SELECT (CTAS) e CREATE DATABASE AS SELECT (CDAS). A instrução CTAS sincroniza o esquema e os dados de uma única tabela, enquanto a CDAS sincroniza um banco de dados inteiro ou várias tabelas dentro do mesmo banco de dados. Para mais informações, consulte Usar instruções CTAS e CDAS no Realtime Compute for Apache Flink para sincronizar dados de um banco de dados MySQL para o StarRocks.
Sintaxe
CREATE TABLE USER_RESULT(
name VARCHAR,
score BIGINT
) WITH (
'connector' = 'starrocks',
'jdbc-url'='jdbc:mysql://fe1_ip:query_port,fe2_ip:query_port,fe3_ip:query_port?xxxxx',
'load-url'='fe1_ip:http_port;fe2_ip:http_port;fe3_ip:http_port',
'database-name' = 'xxx',
'table-name' = 'xxx',
'username' = 'xxx',
'password' = 'xxx'
);
Parâmetros
|
Tipo |
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
|
Geral |
connector |
Especifica o conector a ser usado. |
String |
Sim |
— |
O valor deve ser |
|
jdbc-url |
URL de Java Database Connectivity (JDBC). |
String |
Sim |
— |
Especifique o endereço IP e a porta JDBC do FE no formato |
|
|
database-name |
Nome do banco de dados StarRocks. |
String |
Sim |
— |
— |
|
|
table-name |
Nome da tabela StarRocks. |
String |
Sim |
— |
— |
|
|
username |
Nome de usuário para conexão ao StarRocks. |
String |
Sim |
— |
— |
|
|
password |
Senha para conexão ao StarRocks. |
String |
Sim |
— |
— |
|
|
starrocks.create.table.properties |
Define propriedades para criação automática de tabelas. |
String |
Não |
— |
Especifica propriedades iniciais da tabela, como engine e número de réplicas. Exemplo: 'starrocks.create.table.properties' = 'buckets 8' ou 'starrocks.create.table.properties' = 'replication_num=1'. |
|
|
Específico de source |
scan-url |
URL para varredura de dados. |
String |
Não |
— |
Especifica o endereço IP e a porta http do FE. Formato: Nota
Para especificar vários endereços IP e portas, separe-os por ponto e vírgula (;). |
|
scan.connect.timeout-ms |
Tempo limite para o O conector reporta um erro caso a conexão não seja estabelecida dentro desse prazo. |
String |
Não |
1000 |
Unidade: milissegundos. |
|
|
scan.params.keep-alive-min |
Duração keep-alive para a tarefa de consulta. |
String |
Não |
10 |
— |
|
|
scan.params.query-timeout-s |
Tempo limite para uma tarefa de consulta. Se nenhum resultado for retornado nesse período, o sistema interrompe a tarefa de consulta. |
String |
Não |
600 |
Unidade: segundos. |
|
|
scan.params.mem-limit-byte |
Limite de memória para uma única consulta em um nó BE. |
String |
Não |
1073741824 (1 GB) |
Unidade: bytes. |
|
|
scan.max-retries |
Número máximo de tentativas para uma consulta com falha. O conector reporta um erro se esse limite for excedido. |
String |
Não |
1 |
— |
|
|
Específico de destino |
load-url |
URL para importação de dados. |
String |
Sim |
— |
Especifique os endereços IP e as portas http dos FEs no formato Nota
Para especificar vários endereços IP e portas, separe-os por ponto e vírgula (;). |
|
sink.semantic |
Semântica de entrega para gravações. |
String |
Não |
at-least-once |
Valores válidos:
|
|
|
sink.buffer-flush.max-bytes |
Quantidade máxima de dados a serem armazenados em buffer antes do flush. |
String |
Não |
94371840 (90 MB) |
Intervalo válido: 64 MB a 10 GB. |
|
|
sink.buffer-flush.max-rows |
Número máximo de linhas a serem armazenadas em buffer antes do flush. |
String |
Não |
500000 |
Intervalo válido: 64.000 a 5.000.000. |
|
|
sink.buffer-flush.interval-ms |
Intervalo de flush do buffer. |
String |
Não |
300000 |
Intervalo válido: 1.000 ms a 3.600.000 ms. |
|
|
sink.max-retries |
Número máximo de tentativas para gravações com falha. |
String |
Não |
3 |
Intervalo válido: 0 a 10. |
|
|
sink.connect.timeout-ms |
Tempo limite para conexão ao StarRocks. |
String |
Não |
1000 |
Intervalo válido: 100 a 60.000. Unidade: milissegundos. |
|
|
sink.properties.* |
Propriedades adicionais do Stream Load para o destino. |
String |
Não |
— |
Esses parâmetros controlam o comportamento do Stream Load. Por exemplo, |
|
|
Específico de dimensão |
lookup.cache.enabled |
Defina se o cache deve ser ativado para a tabela de dimensão. |
Boolean |
Não |
true |
Valores válidos:
Importante
|
Mapeamento de tipos de dados
|
Tipo de dados StarRocks |
Tipo de dados Flink |
|
NULL |
NULL |
|
BOOLEAN |
BOOLEAN |
|
TINYINT |
TINYINT |
|
SMALLINT |
SMALLINT |
|
INT |
INT |
|
BIGINT |
BIGINT |
|
BIGINT UNSIGNED Nota
Requer o mecanismo Realtime Compute for Apache Flink VVR 8.0.10 ou posterior. |
DECIMAL(20,0) |
|
LARGEINT |
DECIMAL(20,0) |
|
FLOAT |
FLOAT |
|
DOUBLE |
DOUBLE |
|
DATE |
DATE |
|
DATETIME |
TIMESTAMP |
|
DECIMAL |
DECIMAL |
|
DECIMALV2 |
DECIMAL |
|
DECIMAL32 |
DECIMAL |
|
DECIMAL64 |
DECIMAL |
|
DECIMAL128 |
DECIMAL |
|
CHAR(m) Nota
|
CHAR(n) |
|
VARCHAR(m) Nota
|
CHAR(n) |
|
VARCHAR |
STRING |
|
VARBINARY Nota
Requer o mecanismo Realtime Compute for Apache Flink VVR 8.0.10 ou posterior. |
VARBINARY |
Exemplo de código
CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_source` (
`runoob_id` BIGINT NOT NULL,
`runoob_title` STRING NOT NULL,
`runoob_author` STRING NOT NULL,
`submission_date` DATE NULL
) WITH (
'connector' = 'starrocks',
'jdbc-url' = 'jdbc:mysql://ip:9030',
'scan-url' = 'ip:18030',
'database-name' = 'db_name',
'table-name' = 'table_name',
'password' = 'xxxxxxx',
'username' = 'xxxxx'
);
CREATE TEMPORARY TABLE IF NOT EXISTS `runoob_tbl_sink` (
`runoob_id` BIGINT NOT NULL,
`runoob_title` STRING NOT NULL,
`runoob_author` STRING NOT NULL,
`submission_date` DATE NULL
PRIMARY KEY(`runoob_id`)
NOT ENFORCED
) WITH (
'jdbc-url' = 'jdbc:mysql://ip:9030',
'connector' = 'starrocks',
'load-url' = 'ip:18030',
'database-name' = 'db_name',
'table-name' = 'table_name',
'password' = 'xxxxxxx',
'username' = 'xxxx',
'sink.buffer-flush.interval-ms' = '5000'
);
INSERT INTO runoob_tbl_sink SELECT * FROM runoob_tbl_source;
O StarRocks permite que uma coluna de chave primária seja NULLABLE. No entanto, o Flink não suporta uma chave primária que contenha uma coluna anulável. O modelo de consistência de dados do Flink exige que uma chave primária seja única e não anulável. Caso contrário, o Flink lança o erro Invalid primary key. Column 'xxx' is nullable. Para mais informações, consulte Erro "Invalid primary key. Column 'xxx' is nullable.".
Ingestão de dados
Use o conector StarRocks Pipeline para gravar registros de dados e alterações de esquema de fontes de dados upstream em um banco de dados StarRocks externo. O conector do StarRocks suporta tanto a edição comunitária quanto o EMR Serverless StarRocks totalmente gerenciado da Alibaba Cloud.
Recursos
-
Criação automática de bancos de dados e tabelas.
Se um banco de dados ou tabela upstream não existir na instância StarRocks downstream, o conector o criará automaticamente. Use o parâmetro
table.create.properties.*para configure opções de criação automática de tabelas. -
Sincronização de alterações de esquema.
O conector do StarRocks aplica automaticamente eventos CreateTableEvent, AddColumnEvent e DropColumnEvent ao banco de dados downstream.
O VVR 11.1 e posteriores suportam alterações compatíveis de tipos de colunas. Para mais informações, consulte ALTER TABLE | StarRocks.
Notas de uso
-
Cada tabela sincronizada deve ter uma chave primária. Para tabelas sem chave primária, especifique uma no bloco
transformpara gravar dados downstream. Por exemplo:transform: - source-table: ... primary-keys: id, ... Para tabelas criadas automaticamente, a chave de bucket é igual à chave primária, e a tabela não pode ter uma chave de partição.
Ao sincronizar alterações de esquema, novas colunas só podem ser anexadas ao final das colunas existentes. No modo padrão de alteração de esquema Lenient, inserções em outras posições são movidas automaticamente para o final.
Se você usar uma versão do StarRocks anterior à 2.5.7, especifique explicitamente o número de buckets com o parâmetro
table.create.num-buckets. O StarRocks 2.5.7 e posteriores podem determinar automaticamente um número apropriado de buckets.Caso utilize o StarRocks 3.2 ou posterior, recomendamos ative a opção
table.create.properties.fast_schema_evolutionpara acelerar alterações de esquema.-
Problemas de streaming podem ocorrer ao usar CDC yaml para ingestão de dados no EMR Serverless StarRocks. Utilize uma das seguintes soluções alternativas:
Use o conector Flink SQL StarRocks e defina o parâmetro
sink.version=V1.Ative o parâmetro FE
emr_internal_redirect.Use um nome de domínio StarRocks Private Zone em vez de um SLB.
Sintaxe
source:
...
sink:
type: starrocks
name: StarRocks Sink
jdbc-url: jdbc:mysql://127.0.0.1:9030
load-url: 127.0.0.1:8030
username: root
password: pass
sink.buffer-flush.interval-ms: 5000 # Set the data flush interval.
Configuração
|
Parâmetro |
Descrição |
Tipo |
Obrigatório |
Padrão |
Observações |
|
|
Especifica o tipo do conector de destino. |
String |
Sim |
— |
Defina como |
|
|
Nome de exibição do destino. |
String |
Não |
— |
— |
|
|
URL JDBC para conexão ao banco de dados. |
String |
Sim |
— |
Suporta múltiplos endereços separados por vírgulas ( |
|
|
URL http de um nó FE para Stream Load. |
String |
Sim |
— |
Suporta múltiplos endereços separados por ponto e vírgula ( |
|
|
Nome de usuário para conexão ao StarRocks. |
String |
Sim |
— |
Este usuário deve ter pelo menos permissões SELECT e INSERT na tabela de destino. Conceda as permissões necessárias com o comando GRANT do StarRocks. |
|
|
Senha para conexão ao StarRocks. |
String |
Sim |
— |
— |
|
|
Semântica de entrega para gravações de dados. |
String |
Não |
at-least-once |
Valores válidos:
|
|
|
Prefixo de rótulo para tarefas de Stream Load. |
String |
Não |
— |
— |
|
|
Tempo limite para estabelecer uma conexão http. |
Integer |
Não |
30000 |
Unidade: milissegundos. O valor deve estar entre 100 e 60000. |
|
|
Tempo limite para aguardar uma resposta 100 Continue do servidor. |
Integer |
Não |
30000 |
Unidade: milissegundos. O valor deve estar entre 3000 e 600000. |
|
|
Tamanho máximo do cache em memória, em bytes, antes de acionar um flush. |
Long |
Não |
157286400 |
Unidade: bytes. O valor deve estar entre 64 MB e 10 GB. Nota
|
|
|
Número máximo de linhas no cache em memória antes de acionar um flush. |
Long |
Não |
500000 |
O valor deve estar entre 64.000 e 5.000.000. |
|
|
Intervalo de tempo entre flushes para o buffer de cada tabela. |
Long |
Não |
300000 |
Unidade: milissegundos. Nota
Para tarefas que sincronizam pequenas quantidades de dados, reduza este valor para evitar longos atrasos antes que os dados sejam persistidos. |
|
|
Número máximo de tentativas. |
Long |
Não |
3 |
O valor deve estar entre 0 e 1000. |
|
|
Frequência com que o conector verifica se deve realizar o flush do buffer. |
Long |
Não |
50 |
Unidade: milissegundos. |
|
|
Número de threads usadas para Stream Load. |
Integer |
Não |
2 |
— |
|
|
Defina se deve ser usada a interface de transação Stream Load para ingestão de dados. |
Boolean |
Não |
true |
Esta opção só tem efeito se o banco de dados oferecer suporte a ela. |
|
|
Propriedades adicionais para o destino. |
String |
Não |
— |
Para propriedades suportadas, consulte STREAM LOAD. |
|
|
Número de buckets para tabelas criadas automaticamente. |
Integer |
Não |
— |
|
|
|
Propriedades adicionais para criação automática de tabelas. |
String |
Não |
— |
Por exemplo, passe |
|
|
Tempo limite para operações de alteração de esquema. |
Duration |
Não |
30 min |
Deve ser um número inteiro de segundos. Nota
Se uma operação de alteração de esquema exceder esse limite, a tarefa falhará. |
|
|
Número de bytes a serem alocados para cada caractere Unicode. |
Integer |
Não |
3 |
No CDC, o comprimento de um tipo VARCHAR é medido em caracteres, enquanto no StarRocks, o comprimento de um tipo VARCHAR é medido em bytes. Na maioria dos casos, um caractere Unicode não excede 3 bytes após codificação UTF-8. No entanto, alguns caracteres raros e símbolos emoji podem ocupar 4 ou mais bytes. |
Reutilizar um catálogo integrado
O VVR 11.5 e posteriores permitem referenciar um catálogo StarRocks integrado criado na página Data Management diretamente em uma tarefa de ingestão de dados Flink CDC. Isso simplifica a configuração ao reduzir o número de propriedades que precisam ser definidas manualmente.
sink:
type: starrocks
using.built-in-catalog: starrocks_catalog
As tarefas de ingestão de dados podem reutilizar automaticamente as seguintes opções do catálogo StarRocks:
jdbc-url
http-url
username
password
table.num-buckets
Para substituir esses valores, defina explicitamente as opções yaml correspondentes, que terão precedência.
Mapeamento de tipos
O StarRocks não suporta todos os tipos CDC yaml. Gravar um tipo não suportado no destino causa falha na tarefa. Use a função integrada CAST em uma transformação para converter dados não suportados, ou use uma instrução de projeção para removê-los da tabela de resultados. Para mais informações, consulte Desenvolver uma tarefa de ingestão de dados Flink CDC.
|
Tipo CDC |
Tipo StarRocks |
Observações |
|
TINYINT |
TINYINT |
— |
|
SMALLINT |
SMALLINT |
|
|
INT |
INT |
|
|
BIGINT |
BIGINT |
|
|
FLOAT |
FLOAT |
|
|
DOUBLE |
DOUBLE |
|
|
BOOLEAN |
BOOLEAN |
|
|
DATE |
DATE |
|
|
TIMESTAMP |
DATETIME |
|
|
TIMESTAMP_LTZ |
DATETIME |
|
|
DECIMAL(p, s) |
DECIMAL(p, s) |
Como o StarRocks não suporta DECIMAL para chave primária, o conector converte automaticamente uma coluna de chave primária DECIMAL upstream para VARCHAR no esquema StarRocks sincronizado. |
|
CHAR(n) (n <= 85) |
CHAR(n × 3) |
O CDC mede o comprimento em caracteres, enquanto o StarRocks usa bytes. O conector multiplica o comprimento por 3 para considerar caracteres UTF-8 multibyte. Nota
O comprimento máximo do tipo CHAR do StarRocks é 255. Portanto, apenas tipos CHAR do CDC com comprimento de até 85 são mapeados para o tipo CHAR do StarRocks. Nota
Defina o parâmetro |
|
CHAR(n) (n > 85) |
VARCHAR(n × 3) |
O CDC mede o comprimento em caracteres, enquanto o StarRocks usa bytes. O conector multiplica o comprimento por 3 para considerar caracteres UTF-8 multibyte. Nota
O CDC mede o comprimento em caracteres, enquanto o StarRocks usa bytes. O conector multiplica o comprimento por 3. Como o resultado excede o limite de 255 bytes para o tipo CHAR do StarRocks, ele é mapeado para VARCHAR. Nota
Defina o parâmetro |
|
VARCHAR(n) |
VARCHAR(n × 3) |
O CDC mede o comprimento em caracteres, enquanto o StarRocks usa bytes. O conector multiplica o comprimento por 3 para considerar caracteres UTF-8 multibyte. Nota
Defina o parâmetro |
|
BINARY(n) |
BINARY(n+2) |
Dois bytes de preenchimento são adicionados para garantir a integridade dos dados. |
|
VARBINARY(n) |
VARBINARY(n+1) |
Um byte de preenchimento é adicionado para garantir a integridade dos dados. |
Alteração de esquema
Como destino de ingestão de dados, o StarRocks suporta os seguintes eventos de alteração de esquema:
-
CREATE TABLE EVENT
NotaSe a tabela StarRocks downstream já existir, o conector não tentará crie novamente. Certifique-se de que o esquema da tabela downstream seja compatível com o esquema upstream.
-
ADD COLUMN EVENT
NotaO StarRocks exige que as colunas de chave primária apareçam primeiro em uma tabela. Quaisquer novas colunas devem ser adicionadas após elas.
-
ALTER COLUMN TYPE EVENT
NotaPara caminhos de alteração de esquema suportados, consulte a documentação oficial do StarRocks.
DROP COLUMN EVENT
TRUNCATE TABLE EVENT
DROP TABLE EVENT