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:
Realtime Compute for Apache Flink com Ververica Runtime (VVR) 8.0.10 ou superior
Uma instância do ApsaraDB for SelectDB. Consulte Criar uma instância.
Uma lista de permissões de endereços IP configurada na instância. Consulte Configurar uma lista de permissões.
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:
Baixe o pacote JAR no Maven Central (versões do Flink 1.15–1.17).
Faça upload do JAR para o console de desenvolvimento do Realtime Compute for Apache Flink. Consulte Gerenciar conectores personalizados.
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 |
|
|
Sim |
— |
Fixo como |
|
|
Sim |
— |
Endpoint HTTP da instância do SelectDB: |
|
|
Não |
— |
String de conexão Java Database Connectivity (JDBC) para consultas em tabelas de dimensão e metadados: |
|
|
Sim |
— |
Tabela de destino no formato |
|
|
Sim |
— |
Nome de usuário do banco de dados. Redefina a senha no canto superior direito da página Instance Details, se necessário. |
|
|
Sim |
— |
Senha do usuário do banco de dados. |
|
|
Não |
|
Número de tentativas para solicitações com falha. |
|
|
Não |
|
Tempo limite de conexão. |
|
|
Não |
|
Tempo limite de leitura. |
Tabela source
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Não |
|
Tempo limite da consulta (6 horas por padrão). |
|
|
Não |
|
Número de tablets por partição. Valores menores aumentam o paralelismo do Flink, mas exercem mais pressão sobre o banco de dados. |
|
|
Não |
|
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. |
|
|
Não |
|
Limite de memória por consulta em bytes (8 GB por padrão). |
|
|
Não |
|
Nenhuma configuração necessária — ativar Direct Cluster Connection no console do SelectDB habilita o Arrow Flight SQL automaticamente. |
|
|
Não |
— |
Porta Arrow Flight SQL ( |
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 |
|
|
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. |
|
|
Não |
— |
Parâmetros de importação do Stream Load passados diretamente para a API Stream Load do SelectDB. Veja exemplos abaixo. |
|
|
Não |
|
Propaga operações DELETE. Requer que a tabela Doris tenha exclusão em lote ativada e funciona apenas com o modelo Unique. |
|
|
Não |
|
Ativa o two-phase commit (2PC) para semântica exactly-once. Consulte Explicit Transaction Operations. |
|
|
Não |
|
Tamanho do buffer de cache de gravação em bytes. Mantenha o valor padrão. |
|
|
Não |
|
Quantidade de buffers de cache de gravação. Mantenha o valor padrão. |
|
|
Não |
|
Máximo de tentativas após falha na confirmação. |
|
|
Não |
|
Alterna para o modo de gravação em batch. O flush é controlado pelos três parâmetros |
|
|
Não |
|
Tamanho da fila de cache no modo batch. |
|
|
Não |
|
Máximo de linhas por flush no modo batch. |
|
|
Não |
|
Máximo de bytes por flush no modo batch. |
|
|
Não |
|
Intervalo de flush no modo batch. |
|
|
Não |
|
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 |
|
|
Não |
|
Máximo de linhas no cache de lookup. |
|
|
Não |
|
Tempo de vida (TTL) da entrada de cache. |
|
|
Não |
|
Tentativas após falha na consulta de lookup. |
|
|
Não |
|
Ativa lookup assíncrono. |
|
|
Não |
|
Tamanho máximo do lote por consulta no modo de lookup assíncrono. |
|
|
Não |
|
Tamanho da fila de buffer intermediário no modo de lookup assíncrono. |
|
|
Não |
|
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 |
|
|
Sim |
— |
Fixo como |
|
|
Não |
— |
Nome descritivo para o sink. |
|
|
Sim |
— |
Endpoint HTTP: |
|
|
Não |
— |
String de conexão JDBC. Exemplo: |
|
|
Sim |
— |
Nome de usuário do banco de dados. |
|
|
Sim |
— |
Senha do usuário do banco de dados. |
|
|
Não |
|
O modo batch está ativado por padrão em jobs de ingestão de dados. O flush é controlado pelos três parâmetros |
|
|
Não |
|
Tamanho da fila de cache. |
|
|
Não |
|
Máximo de linhas por flush. |
|
|
Não |
|
Máximo de bytes por flush. |
|
|
Não |
|
Intervalo de flush. Mínimo: |
|
|
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 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
O SelectDB armazena strings em UTF-8. Caracteres em inglês ocupam 1 byte; caracteres chineses ocupam 3 bytes. O comprimento máximo de |
|
|
|
Aplica-se o mesmo multiplicador UTF-8. O comprimento máximo de |
|
|
|
|
|
|
|
|
|
|
|