Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:StarRocks

Última atualização: Jul 07, 2026

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 starrocks.

jdbc-url

URL de Java Database Connectivity (JDBC).

String

Sim

Especifique o endereço IP e a porta JDBC do FE no formato jdbc:mysql://ip:port.

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: fe_ip:http_port;fe_ip:http_port.

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 flink-connector-starrocks se conectar ao StarRocks.

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 fe_ip:http_port;fe_ip:http_port.

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:

  • at-least-once (padrão): Garante que os dados sejam entregues pelo menos uma vez.

  • exactly-once: Garante que os dados sejam entregues exatamente uma vez.

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, sink.properties.format especifica o formato dos dados importados, como CSV. Para mais parâmetros, consulte Stream Load.

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:

  • true: Ative o cache. Após a primeira leitura dos dados da tabela, o sistema os armazena em cache na memória. Solicitações subsequentes usam os dados em cache dentro do período de validade para reduzir a sobrecarga de I/O.

  • false: Desativa o cache. Cada consulta acessa diretamente a fonte de dados.

Importante
  • Este recurso requer o mecanismo Realtime Compute for Apache Flink VVR 11.1 ou posterior.

  • Recomendamos desativar este recurso nos seguintes cenários:

    • Os dados da tabela de dimensão são atualizados frequentemente e dados em tempo real são necessários.

    • A tabela contém uma grande quantidade de dados, apresentando risco de estouro de memória.

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
  • O VVR 8.0.10 estende automaticamente o comprimento CHAR em três vezes (m=n*3, onde n<=85) para acomodar diferenças de codificação entre MySQL e StarRocks.

  • O VVR 8.0.11 e posteriores estendem automaticamente o comprimento CHAR em quatro vezes (m=n*4, onde n<=63) para acomodar diferenças de codificação entre MySQL e StarRocks.

  • O comprimento máximo para o tipo CHAR do StarRocks é 255. Portanto, o Flink mapeia um tipo CHAR para um tipo CHAR do StarRocks somente se seu comprimento estendido automaticamente não exceder 255.

CHAR(n)

VARCHAR(m)

Nota
  • O VVR 8.0.10 estende automaticamente o comprimento VARCHAR em três vezes (m=n*3, onde n>85) para acomodar diferenças de codificação entre MySQL e StarRocks.

  • O VVR 8.0.11 e posteriores estendem automaticamente o comprimento VARCHAR em quatro vezes (m=n*4, onde n>63) para acomodar diferenças de codificação entre MySQL e StarRocks.

  • O comprimento máximo para o tipo CHAR do StarRocks é 255. Portanto, se o comprimento estendido automaticamente de um tipo CHAR do Flink exceder 255, o Flink mapeia o tipo para o tipo VARCHAR do StarRocks.

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;
Nota

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 transform para 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_evolution para 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

type

Especifica o tipo do conector de destino.

String

Sim

Defina como starrocks.

name

Nome de exibição do destino.

String

Não

jdbc-url

URL JDBC para conexão ao banco de dados.

String

Sim

Suporta múltiplos endereços separados por vírgulas (,). Exemplo: jdbc:mysql://fe_host1:fe_query_port1,fe_host2:fe_query_port2,fe_host3:fe_query_port3.

load-url

URL http de um nó FE para Stream Load.

String

Sim

Suporta múltiplos endereços separados por ponto e vírgula (;). Exemplo: fe_host1:fe_http_port1;fe_host2:fe_http_port2.

username

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.

password

Senha para conexão ao StarRocks.

String

Sim

sink.semantic

Semântica de entrega para gravações de dados.

String

Não

at-least-once

Valores válidos:

  • at-least-once (padrão): Garante que os dados sejam entregues pelo menos uma vez.

  • exactly-once: Garante que os dados sejam entregues exatamente uma vez.

sink.label-prefix

Prefixo de rótulo para tarefas de Stream Load.

String

Não

sink.connect.timeout-ms

Tempo limite para estabelecer uma conexão http.

Integer

Não

30000

Unidade: milissegundos. O valor deve estar entre 100 e 60000.

sink.wait-for-continue.timeout-ms

Tempo limite para aguardar uma resposta 100 Continue do servidor.

Integer

Não

30000

Unidade: milissegundos. O valor deve estar entre 3000 e 600000.

sink.buffer-flush.max-bytes

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
  • Este tamanho de cache é compartilhado por todas as tabelas. Quando o buffer estiver cheio, o conector selecione algumas tabelas para realizar o flush.

  • Definir um valor maior pode melhorar o throughput, mas pode aumentar a latência de ingestão.

sink.buffer-flush.max-rows

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.

sink.buffer-flush.interval-ms

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.

sink.max-retries

Número máximo de tentativas.

Long

Não

3

O valor deve estar entre 0 e 1000.

sink.scan-frequency.ms

Frequência com que o conector verifica se deve realizar o flush do buffer.

Long

Não

50

Unidade: milissegundos.

sink.io.thread-count

Número de threads usadas para Stream Load.

Integer

Não

2

sink.at-least-once.use-transaction-stream-load

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.

sink.properties.*

Propriedades adicionais para o destino.

String

Não

Para propriedades suportadas, consulte STREAM LOAD.

table.create.num-buckets

Número de buckets para tabelas criadas automaticamente.

Integer

Não

table.create.properties.*

Propriedades adicionais para criação automática de tabelas.

String

Não

Por exemplo, passe 'table.create.properties.fast_schema_evolution' = 'true' para ative alteração rápida de esquema. Para detalhes, consulte a documentação do StarRocks.

table.schema-change.timeout

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á.

unicode-char.max-bytes

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

Nota

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 unicode-char.max-bytes para alocar mais bytes para cada caractere Unicode.

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 unicode-char.max-bytes para alocar mais bytes para cada caractere Unicode.

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 unicode-char.max-bytes para alocar mais bytes para cada caractere Unicode.

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

    Nota

    Se 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

    Nota

    O StarRocks exige que as colunas de chave primária apareçam primeiro em uma tabela. Quaisquer novas colunas devem ser adicionadas após elas.

  • DROP COLUMN EVENT

  • TRUNCATE TABLE EVENT

  • DROP TABLE EVENT