Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:SelectDB

Última atualização: Jun 27, 2026

O conector do SelectDB integra o Realtime Compute for Apache Flink ao ApsaraDB for SelectDB, um data warehouse em tempo real totalmente gerenciado e compatível com Apache Doris no Alibaba Cloud. Utilize-o para construir pipelines em tempo real que leem, gravam ou consultam dados no SelectDB, além de executar sincronização completa de banco de dados em jobs de ingestão de dados baseados em YAML.

Capacidades suportadas:

Categoria

Detalhes

Tipos de tabela

Tabela source, tabela sink, tabela de dimensão, sink de ingestão de dados

Modo de execução

Stream e batch

Formato de dados

JSON e CSV

Tipo de API

DataStream, SQL e jobs YAML de ingestão de dados

Suporte a Update/Delete

Sim

Métricas de monitoramento

Nenhuma

Principais recursos:

  • Sincronização de dados de banco de dados completo

  • Semântica exactly-once via two-phase commit (2PC) — sem registros duplicados ou perdidos

  • Compatível com Apache Doris 1.0 e versões posteriores

Pré-requisitos

Antes de começar, verifique se você possui:

Configurar o conector

O conector do SelectDB está integrado ao VVR 11.1 e versões posteriores — nenhuma instalação manual é necessária.

Para as versões do VVR entre 8.0.10 e 11.0, instale o conector manualmente:

  1. Baixe o pacote JAR no Maven Central (versões do Flink 1.15–1.17).

  2. Faça upload do JAR para o console de desenvolvimento do Realtime Compute for Apache Flink. Consulte Gerenciar conectores personalizados.

  3. Referencie o conector no seu job SQL usando 'connector' = 'doris'.

SQL

Sintaxe

Os três tipos de tabela — source, sink e dimensão — compartilham a mesma sintaxe DDL. Especifique a função da tabela por meio dos parâmetros incluídos.

Para usar o SelectDB como tabela source, ative primeiro a conexão direta ao cluster. No console do ApsaraDB for SelectDB, acesse Instance Details > Network Information e clique em Enable Direct Cluster Connection . Isso ativa o protocolo Arrow Flight SQL para leituras paralelas de alto throughput.
CREATE TABLE selectdb_source (
  order_id      BIGINT,
  user_id       BIGINT,
  total_amount  DECIMAL(10, 2),
  order_status  TINYINT,
  create_time   TIMESTAMP(3),
  product_name  STRING
) WITH (
  'connector'        = 'doris',
  'fenodes'          = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
  'table.identifier' = 'shop_db.orders',
  'username'         = 'admin',
  'password'         = '****'
);

Parâmetros

Geral

Parâmetro

Obrigatório

Padrão

Descrição

connector

Sim

Fixo como doris.

fenodes

Sim

Endpoint HTTP da instância do SelectDB: <VPC Address or Public Address>:<HTTP Protocol Port>. Obtenha ambos em Instance Details > Network Information no console do SelectDB. Exemplo: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080.

jdbc-url

Não

String de conexão Java Database Connectivity (JDBC) para consultas em tabelas de dimensão e metadados: jdbc:mysql://<VPC Address or Public Address>:<MySQL Protocol Port>. Exemplo: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030.

table.identifier

Sim

Tabela de destino no formato <database>.<table>. Exemplo: db.tbl.

username

Sim

Nome de usuário do banco de dados. Redefina a senha no canto superior direito da página Instance Details, se necessário.

password

Sim

Senha do usuário do banco de dados.

doris.request.retries

Não

3

Número de tentativas para solicitações com falha.

doris.request.connect.timeout

Não

30s

Tempo limite de conexão.

doris.request.read.timeout

Não

30s

Tempo limite de leitura.

Tabela source

Parâmetro

Obrigatório

Padrão

Descrição

doris.request.query.timeout

Não

21600s

Tempo limite da consulta (6 horas por padrão).

doris.request.tablet.size

Não

1

Número de tablets por partição. Valores menores aumentam o paralelismo do Flink, mas exercem mais pressão sobre o banco de dados.

doris.batch.size

Não

4064

Máximo de linhas lidas de um nó Backend (BE) por solicitação. Aumente para reduzir a sobrecarga de conexão e a latência de rede.

doris.exec.mem.limit

Não

8192mb

Limite de memória por consulta em bytes (8 GB por padrão).

source.use-flight-sql

Não

false

Nenhuma configuração necessária — ativar Direct Cluster Connection no console do SelectDB habilita o Arrow Flight SQL automaticamente.

source.flight-sql-port

Não

Porta Arrow Flight SQL (arrow_flight_sql_port) do nó Frontend (FE).

Tabela sink

O modo de gravação afeta as garantias de entrega e o comportamento de flush. Escolha com base nos seus requisitos de consistência:

Gravação em streaming

Gravação em batch

Condição de gatilho

Segue os intervalos de checkpoint do Flink

Flush periódico por volume de dados ou limiar de tempo

Garantia de entrega

Exactly-once (via 2PC)

At-least-once; obtenha idempotência com o modelo Unique

Latência

Limitada pelo intervalo de checkpoint

Flexível, independente de checkpoints

Tolerância a falhas

Recuperação completa de estado do Flink

Depende da desduplicação do modelo Unique

Parâmetro

Obrigatório

Padrão

Descrição

sink.label-prefix

Não

Prefixo de rótulo para importações via Stream Load. Deve ser globalmente único em todos os jobs — o mesmo rótulo só pode ser confirmado uma vez. Necessário para garantir a semântica exactly-once entre reinicializações de jobs.

sink.properties.*

Não

Parâmetros de importação do Stream Load passados diretamente para a API Stream Load do SelectDB. Veja exemplos abaixo.

sink.enable-delete

Não

true

Propaga operações DELETE. Requer que a tabela Doris tenha exclusão em lote ativada e funciona apenas com o modelo Unique.

sink.enable-2pc

Não

true

Ativa o two-phase commit (2PC) para semântica exactly-once. Consulte Explicit Transaction Operations.

sink.buffer-size

Não

1 MB

Tamanho do buffer de cache de gravação em bytes. Mantenha o valor padrão.

sink.buffer-count

Não

3

Quantidade de buffers de cache de gravação. Mantenha o valor padrão.

sink.max-retries

Não

3

Máximo de tentativas após falha na confirmação.

sink.enable.batch-mode

Não

false

Alterna para o modo de gravação em batch. O flush é controlado pelos três parâmetros sink.buffer-flush.* abaixo, em vez de checkpoints. A garantia exactly-once não se aplica; use o modelo Unique para idempotência.

sink.flush.queue-size

Não

2

Tamanho da fila de cache no modo batch.

sink.buffer-flush.max-rows

Não

500000

Máximo de linhas por flush no modo batch.

sink.buffer-flush.max-bytes

Não

100 MB

Máximo de bytes por flush no modo batch.

sink.buffer-flush.interval

Não

10s

Intervalo de flush no modo batch.

sink.ignore.update-before

Não

true

Ignora eventos update-before do Flink CDC.

**Exemplos de sink.properties.*:**

Formato CSV:

'sink.properties.column_separator' = ','
-- If values may contain commas, use a non-printable separator:
-- 'sink.properties.column_separator' = '\x01'

Formato JSON:

'sink.properties.format'            = 'json',
'sink.properties.read_json_by_line' = 'true'
-- Alternatively: 'sink.properties.strip_outer_array' = 'true'

Tabela de dimensão

Parâmetro

Obrigatório

Padrão

Descrição

lookup.cache.max-rows

Não

-1

Máximo de linhas no cache de lookup. -1 desativa o cache.

lookup.cache.ttl

Não

10s

Tempo de vida (TTL) da entrada de cache.

lookup.max-retries

Não

1

Tentativas após falha na consulta de lookup.

lookup.jdbc.async

Não

false

Ativa lookup assíncrono.

lookup.jdbc.read.batch.size

Não

128

Tamanho máximo do lote por consulta no modo de lookup assíncrono.

lookup.jdbc.read.batch.queue-size

Não

256

Tamanho da fila de buffer intermediário no modo de lookup assíncrono.

lookup.jdbc.read.thread-size

Não

3

Threads de lookup JDBC por tarefa no modo de lookup assíncrono.

Exemplos

Tabela source

CREATE TEMPORARY TABLE selectdb_source (
  order_id      BIGINT,
  user_id       BIGINT,
  total_amount  DECIMAL(10, 2),
  order_status  TINYINT,
  create_time   TIMESTAMP(3),
  product_name  STRING
) WITH (
  'connector'        = 'doris',
  'fenodes'          = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
  'table.identifier' = 'shop_db.orders',
  'username'         = 'admin',
  'password'         = '****'
);

Tabela sink

CREATE TEMPORARY TABLE selectdb_sink (
  order_id      BIGINT,
  user_id       BIGINT,
  total_amount  DECIMAL(10, 2),
  order_status  TINYINT,
  create_time   TIMESTAMP(3),
  product_name  STRING
) WITH (
  'connector'        = 'doris',
  'fenodes'          = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
  'table.identifier' = 'shop_db.orders',
  'username'         = 'admin',
  'password'         = '****',
  'sink.label-prefix' = 'flink_orders'  -- Must be globally unique across jobs
);

Tabela de dimensão

O SelectDB atua como uma tabela de dimensão de lookup unida a uma tabela de fatos em streaming.

-- Fact table from Kafka
CREATE TEMPORARY TABLE fact_table (
  `id`           BIGINT,
  `name`         STRING,
  `city`         STRING,
  `process_time` AS proctime()
) WITH (
  'connector' = 'kafka',
  ...
);

-- Dimension table from SelectDB
CREATE TEMPORARY TABLE dim_city (
  `city`     STRING,
  `level`    INT,
  `province` STRING,
  `country`  STRING
) WITH (
  'connector'        = 'doris',
  'fenodes'          = 'selectdb-cn-*******.selectdbfe.rds.aliyuncs.com:8080',
  'jdbc-url'         = 'jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030',
  'table.identifier' = 'dim.dim_city',
  'username'         = 'admin',
  'password'         = '****'
);

-- Temporal join
SELECT a.id, a.name, a.city, c.province, c.country, c.level
FROM fact_table a
LEFT JOIN dim_city FOR SYSTEM_TIME AS OF a.process_time AS c
ON a.city = c.city;

Ingestão de dados

Utilize o conector do SelectDB como sink em jobs de ingestão de dados baseados em YAML para sincronização completa de banco de dados.

Sintaxe

source:
  type: <source-type>

sink:
  type: doris
  name: Doris Sink
  fenodes: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080
  username: root
  password: ""

Parâmetros

Parâmetro

Obrigatório

Padrão

Descrição

type

Sim

Fixo como doris.

name

Não

Nome descritivo para o sink.

fenodes

Sim

Endpoint HTTP: <VPC Address or Public Address>:<HTTP Protocol Port>. Obtenha ambos em Instance Details > Network Information no console do SelectDB. Exemplo: selectdb-cn-****.selectdbfe.rds.aliyuncs.com:8080.

jdbc-url

Não

String de conexão JDBC. Exemplo: jdbc:mysql://selectdb-cn-***.selectdbfe.rds.aliyuncs.com:9030.

username

Sim

Nome de usuário do banco de dados.

password

Sim

Senha do usuário do banco de dados.

sink.enable.batch-mode

Não

true

O modo batch está ativado por padrão em jobs de ingestão de dados. O flush é controlado pelos três parâmetros sink.buffer-flush.*. A garantia exactly-once não se aplica; use o modelo Unique para idempotência.

sink.flush.queue-size

Não

2

Tamanho da fila de cache.

sink.buffer-flush.max-rows

Não

500000

Máximo de linhas por flush.

sink.buffer-flush.max-bytes

Não

100 MB

Máximo de bytes por flush.

sink.buffer-flush.interval

Não

10s

Intervalo de flush. Mínimo: 1s.

sink.properties.*

Não

Parâmetros de importação do Stream Load.

**Exemplos de sink.properties.*:**

Formato CSV:

sink.properties.column_separator: ','
# If values may contain commas, use a non-printable separator:
# sink.properties.column_separator: '\x01'

Formato JSON:

sink.properties.format: 'json'
sink.properties.read_json_by_line: 'true'

Mapeamento de tipos

Flink para SelectDB

Tipo Flink CDC

Tipo SelectDB

Observações

TINYINT

TINYINT

SMALLINT

SMALLINT

INT

INT

BIGINT

BIGINT

DECIMAL

DECIMAL

FLOAT

FLOAT

DOUBLE

DOUBLE

BOOLEAN

BOOLEAN

DATE

DATE

TIMESTAMP[(p)]

DATETIME[(p)]

TIMESTAMP_LTZ[(p)]

DATETIME[(p)]

CHAR(n)

CHAR(n*3)

O SelectDB armazena strings em UTF-8. Caracteres em inglês ocupam 1 byte; caracteres chineses ocupam 3 bytes. O comprimento máximo de CHAR é 255; valores maiores são convertidos automaticamente para VARCHAR.

VARCHAR(n)

VARCHAR(n*3)

Aplica-se o mesmo multiplicador UTF-8. O comprimento máximo de VARCHAR é 65533; valores maiores são convertidos automaticamente para STRING.

BINARY(n)

STRING

VARBINARY(n)

STRING

STRING

STRING

Próximos passos