O conector do OceanBase integra o Realtime Compute for Apache Flink ao OceanBase, um banco de dados nativo distribuído para processamento transacional e analítico híbrido (HTAP). Use-o para ler fluxos de captura de dados de alteração (CDC), juntar tabelas de dimensão e gravar resultados no OceanBase.
O conector do OceanBase está em visualização pública.
Informações básicas
O OceanBase é um sistema de gerenciamento de banco de dados nativo distribuído para processamento transacional e analítico híbrido (HTAP). Para mais informações, consulte o site do OceanBase. Para reduzir custos de refatoração dos sistemas de negócios durante a migração de um banco de dados MySQL ou Oracle, o OceanBase oferece dois modos de compatibilidade: modo Oracle e modo MySQL. Os tipos de dados, recursos SQL e visualizações internas nesses modos são consistentes com os do MySQL ou Oracle.
Recursos suportados
|
Categoria |
Detalhes |
|
Tipos de tabela |
Tabelas de source, dimensão e destino |
|
Modos de execução |
Modo streaming e modo batch |
|
Formato de dados |
Não aplicável |
|
Métricas específicas de monitoramento |
Nenhuma |
|
API |
SQL |
|
Atualizações e exclusões em tabelas de destino |
Sim |
Pré-requisitos
Antes de começar, certifique-se de que:
O banco de dados e as tabelas de destino existem no OceanBase
Uma lista de permissões de endereços IP está configurada — consulte Configure a whitelist group
(Para tabelas de source CDC) O OceanBase Binlog service está ativado — consulte Binlog-related operations
(Para tabelas de destino com importação bypass) A porta de importação bypass está ativada — consulte Bypass import
Limitações
Requer Ververica Runtime (VVR) 8.0.1 ou superior.
Garantias semânticas:
Tabelas de source CDC: semântica exactly-once. Não há perda ou duplicação de dados na transição da leitura histórica completa para o Binlog, mesmo após falhas.
Tabelas de destino: semântica at-least-once. Se a tabela de destino tiver chave primária, a idempotência garante a correção dos dados.
Alteração na arquitetura CDC do VVR 11.4.0
A partir do VVR 11.4.0, o conector CDC do OceanBase foi atualizado:
O conector CDC original, baseado no service LogProxy do OceanBase, foi descontinuado e removido.
A captura incremental de logs agora exige o OceanBase Binlog service. O conector CDC do OceanBase oferece melhor compatibilidade de protocolo e estabilidade de conexão com o service Binlog em comparação à conexão direta do conector CDC padrão do MySQL. Não recomendamos conectar o conector CDC padrão do MySQL ao service Binlog do OceanBase para rastreamento de alterações.
O rastreamento incremental de alterações no modo de compatibilidade Oracle não é mais suportado. Para CDC no modo Oracle, entre em contato com o Suporte Técnico Empresarial do OceanBase.
Sintaxe
CREATE TABLE oceanbase_source (
order_id INT,
order_date TIMESTAMP(0),
customer_name STRING,
price DECIMAL(10, 5),
product_id INT,
order_status BOOLEAN,
PRIMARY KEY(order_id) NOT ENFORCED
) WITH (
'connector' = 'oceanbase',
'url' = '<your-jdbc-url>',
'tableName' = '<your-table-name>',
'userName' = '<your-username>',
'password' = '<your-password>'
);
Comportamento de gravação no destino:
Para cada registro recebido, o conector constrói uma instrução SQL com base no esquema da tabela de destino:
Sem chave primária:
INSERT INTOCom chave primária:
UPSERT(dependendo do modo de compatibilidade do banco de dados)
Parâmetros WITH
Parâmetros gerais
Estes parâmetros aplicam-se a todos os tipos de tabela.
|
Parâmetro |
Descrição |
Obrigatório |
Tipo |
Padrão |
|
|
Defina como |
Sim |
STRING |
— |
|
|
Senha do banco de dados. |
Sim |
STRING |
— |
Parâmetros da tabela de source
A partir do VVR 11.4.0, o conector CDC do OceanBase utiliza o service Binlog do OceanBase para captura incremental de logs. O conector baseado em LogProxy foi removido. O CDC no modo de compatibilidade Oracle não é mais suportado a partir do VVR 11.4.0.
| Parâmetro | Descrição | Obrigatório | Tipo | Padrão | Observações |
|---|---|---|---|---|---|
hostname |
Endereço IP ou nome do host do banco de dados OceanBase. Sempre que possível, use um endereço de virtual private cloud (VPC). | Sim | STRING | — | Se o OceanBase e o Realtime Compute for Apache Flink estiverem em VPCs diferentes, configure a conectividade entre VPCs ou use um endpoint público. Consulte Workspace management e Internet access for Flink clusters. |
username |
Nome de usuário do banco de dados OceanBase. | Sim | STRING | — | — |
database-name |
Nome do banco de dados OceanBase. Aceita expressões regulares para leitura de múltiplos bancos de dados. Evite os âncoras ^ e $. |
Sim | STRING | — | O conector concatena database-name e table-name com \\. (VVR 8.0.1+) ou . (versões anteriores) para formar uma regex de caminho completo. Por exemplo, db_.* + tb_.+ torna-se db_.*\\.tb_.+. |
table-name |
Nome da tabela OceanBase. Aceita expressões regulares para leitura de múltiplas tabelas. Evite os âncoras ^ e $. |
Sim | STRING | — | Consulte a observação sobre database-name acima. |
port |
Porta do banco de dados OceanBase. | Não | INTEGER | 3306 |
— |
server-id |
ID numérico para o cliente do banco de dados. Deve ser globalmente único. Aceita um intervalo, como 5400-5408, para atribuir IDs diferentes a leitores concorrentes. |
Não | STRING | Valor aleatório entre 5400 e 6400 | Use um ID diferente para cada job conectado ao mesmo banco de dados. Consulte Server ID usage. |
scan.incremental.snapshot.chunk.size |
Número de linhas por chunk durante a leitura incremental de snapshot. Os dados de cada chunk são armazenados em buffer na memória antes da leitura completa. Chunks menores melhoram a granularidade da recuperação de falhas, mas podem causar erros de falta de memória (OOM) e reduzir o throughput. | Não | INTEGER | 8096 |
Equilibre o tamanho do chunk com os requisitos de memória e throughput. |
scan.snapshot.fetch.size |
Número máximo de registros buscados por pull durante leituras completas de tabela. | Não | INTEGER | 1024 |
— |
scan.startup.mode |
Modo de inicialização para consumo de dados. | Não | STRING | initial |
Valores válidos:
initial
(varre histórico completo e depois lê Binlog),
latest-offset
(apenas cauda do Binlog),
earliest-offset
(Binlog disponível mais antigo),
specific-offset
(definido via parâmetros
scan.startup.specific-offset.*
),
timestamp
(definido via
scan.startup.timestamp-millis
).
Importante
Para os modos |
scan.startup.specific-offset.file |
Nome do arquivo Binlog para o offset inicial. Exemplo: mysql-bin.000003. |
Não | STRING | — | Requer scan.startup.mode=specific-offset. |
scan.startup.specific-offset.pos |
Offset em bytes dentro do arquivo Binlog especificado. | Não | INTEGER | — | Requer scan.startup.mode=specific-offset. |
scan.startup.specific-offset.gtid-set |
Conjunto GTID para o offset inicial. Exemplo: 24DA167-0C0C-11E8-8442-00059A3C7B00:1-19. |
Não | STRING | — | Requer scan.startup.mode=specific-offset. |
scan.startup.timestamp-millis |
Timestamp inicial em milissegundos. O CDC do OceanBase lê o evento inicial de cada arquivo Binlog para localizar o arquivo correspondente a este timestamp. O arquivo Binlog não deve ter sido expurgado. | Não | LONG | — | Requer scan.startup.mode=timestamp. |
server-time-zone |
Fuso horário da sessão usado pelo banco de dados. Controla como os tipos TIMESTAMP são convertidos para STRING. Consulte tipos temporais Debezium. Exemplo: Asia/Shanghai. |
Não | STRING | Fuso horário de execução do job Flink | — |
debezium.min.row.count.to.stream.results |
Limiar de contagem de linhas acima do qual o conector muda da leitura completa (tabela inteira na memória) para leitura em lote (streaming de linhas em lotes). A leitura completa é mais rápida; a leitura em lote evita OOM em tabelas grandes. | Não | INTEGER | 1000 |
— |
connect.timeout |
Tempo máximo de espera antes de tentar novamente uma conexão que atingiu timeout. | Não | DURATION | 30s |
— |
connect.max-retries |
Número máximo de tentativas de reconexão após falha. | Não | INTEGER | 3 |
— |
connection.pool.size |
Tamanho do pool de conexões do banco de dados. Reutilizar conexões reduz o número total de conexões abertas. | Não | INTEGER | 20 |
— |
jdbc.properties.* |
Parâmetros personalizados de conexão URL JDBC. Exemplo: 'jdbc.properties.useSSL' = 'false'. Consulte propriedades de configuração MySQL. |
Não | STRING | — | — |
debezium.* |
Parâmetros personalizados Debezium para leitura de Binlog. Exemplo: 'debezium.event.deserialization.failure.handling.mode' = 'ignore'. |
Não | STRING | — | — |
heartbeat.interval |
Intervalo em que a source emite eventos de heartbeat para avançar o offset do Binlog. Impede que o offset do Binlog expire em tabelas pouco atualizadas. Um offset expirado causa falha no job e exige reinicialização sem estado. | Não | DURATION | 30s |
— |
scan.incremental.snapshot.chunk.key-column |
Coluna usada como chave de chunk para dividir dados durante a fase de snapshot. | Condicional | STRING | — | Obrigatório para tabelas sem chave primária (deve ser NOT NULL). Opcional para tabelas com chave primária (selecione uma coluna da chave primária). |
scan.incremental.close-idle-reader.enabled |
Define se leitores ociosos devem ser fechados após a conclusão da fase de snapshot. | Não | BOOLEAN | false |
VVR 8.0.1+. Também requer execution.checkpointing.checkpoints-after-tasks-finish.enabled=true. |
scan.read-changelog-as-append-only.enabled |
Define se o fluxo de changelog deve ser convertido para um fluxo append-only. Quando true, todos os tipos de mensagem (INSERT, DELETE, UPDATE_BEFORE, UPDATE_AFTER) são convertidos para INSERT. Ative apenas em cenários especiais, como preservar mensagens de exclusão de uma tabela upstream. |
Não | BOOLEAN | false |
VVR 8.0.8+. |
scan.only.deserialize.captured.tables.changelog.enabled |
Define se eventos de alteração devem ser desserializados apenas para as tabelas capturadas durante a fase incremental. Definir como true acelera a leitura do Binlog. |
Não | BOOLEAN | false (VVR 8.x), true (VVR 11.1+) |
VVR 8.0.7+. No VVR 8.0.8 e anteriores, use o nome de parâmetro debezium.scan.only.deserialize.captured.tables.changelog.enable. |
scan.parse.online.schema.changes.enabled |
Define se eventos DDL devem ser analisados para alterações sem bloqueio do ApsaraDB RDS durante a fase incremental. Recurso experimental. Faça um snapshot do job Flink antes de realizar alterações de esquema online sem bloqueio. | Não | BOOLEAN | false |
VVR 11.1+. |
scan.incremental.snapshot.backfill.skip |
Define se o backfill deve ser ignorado durante a fase de snapshot. O backfill aplica-se apenas durante a consulta de snapshot de um único chunk e não cobre toda a fase de leitura completa. Ao ignorar o backfill, a consulta de snapshot de cada chunk lê os dados mais recentes da tabela naquele instante; atualizações ocorridas em um chunk após sua leitura não são mescladas durante a fase de leitura completa e são lidas do Binlog após entrar na fase incremental. Por exemplo, uma atualização no chunk5 ocorrida enquanto o snapshot do chunk5 está sendo feito reflete-se diretamente no snapshot do chunk5; se o chunk5 for atualizado depois que o leitor avançou para o chunk80, a atualização é aplicada posteriormente a partir do Binlog durante a fase incremental. Importante: quando ativado, alterações ocorridas durante ou após a varredura de um chunk ainda são entregues pelo Binlog na fase incremental e podem ser duplicadas; apenas a semântica at-least-once é garantida. Ative isso somente quando o sink downstream suportar gravações idempotentes por chave primária. | Não | BOOLEAN | false |
VVR 11.1+. |
scan.incremental.snapshot.unbounded-chunk-first.enabled |
Define se o chunk ilimitado deve ser distribuído primeiro durante a fase de snapshot. Reduz o risco de OOM quando um TaskManager processa o último chunk. Recurso experimental. Adicione este parâmetro antes do job iniciar pela primeira vez. | Não | BOOLEAN | false |
VVR 11.1+. |
Parâmetros da tabela de dimensão
|
Parâmetro |
Descrição |
Obrigatório |
Tipo |
Padrão |
Observações |
|
|
URL JDBC. Deve incluir o nome do banco de dados MySQL ou o nome do service Oracle. |
Sim |
STRING |
— |
— |
|
|
Nome de usuário do banco de dados. |
Sim |
STRING |
— |
— |
|
|
Política de cache para consultas à tabela de dimensão. |
Não |
STRING |
|
|
|
|
Número máximo de entradas em cache. |
Não |
INTEGER |
|
Obrigatório quando |
|
|
Timeout do cache em milissegundos. O comportamento depende da configuração de |
Não |
LONG |
|
— |
|
|
Duração máxima de nova tentativa. |
Não |
DURATION |
|
— |
Parâmetros da tabela de destino (JDBC)
|
Parâmetro |
Descrição |
Obrigatório |
Tipo |
Padrão |
Observações |
|
|
URL JDBC. Deve incluir o nome do banco de dados MySQL ou o nome do service Oracle. |
Sim |
STRING |
— |
— |
|
|
Nome de usuário do banco de dados. |
Sim |
STRING |
— |
— |
|
|
Nome da tabela de destino. |
Sim |
STRING |
— |
— |
|
|
Modo de gravação. Defina como |
Sim |
STRING |
|
— |
|
|
Modo de compatibilidade do OceanBase. Valores válidos: |
Não |
STRING |
|
Parâmetro específico do OceanBase. |
|
|
Número máximo de novas tentativas de gravação. |
Não |
INTEGER |
|
— |
|
|
Tamanho inicial do pool de conexões. |
Não |
INTEGER |
|
— |
|
|
Número máximo de conexões ativas no pool. |
Não |
INTEGER |
|
— |
|
|
Tempo máximo de espera (ms) por uma conexão do pool. |
Não |
INTEGER |
|
— |
|
|
Número mínimo de conexões ociosas no pool. |
Não |
INTEGER |
|
— |
|
|
Propriedades de conexão JDBC no formato |
Não |
STRING |
— |
— |
|
|
Define se operações de exclusão devem ser ignoradas. |
Não |
BOOLEAN |
|
— |
|
|
Colunas a excluir das atualizações, separadas por vírgula (por exemplo, |
Não |
STRING |
— |
— |
|
|
Chave de partição. Quando definida, os dados são agrupados por esta chave antes da aplicação da |
Não |
STRING |
— |
— |
|
|
Regra de agrupamento no formato |
Não |
STRING |
— |
— |
|
|
Tamanho do buffer de dados (número de registros). |
Não |
INTEGER |
|
— |
|
|
Intervalo de liberação do buffer (ms). Se o buffer não atingir a condição de saída dentro deste intervalo, todos os dados em buffer serão liberados automaticamente. |
Não |
LONG |
|
— |
|
|
Intervalo de nova tentativa (ms). |
Não |
INTEGER |
|
— |
Parâmetros da tabela de destino (importação bypass)
A importação bypass é um método de gravação de alto throughput para carregamento em massa de dados no OceanBase. Disponível no VVR 11.5 e posteriores.
Antes de usar a importação bypass, verifique as seguintes restrições:
Apenas fluxos delimitados: A fonte de dados deve ser um fluxo delimitado. Use o modo batch do Flink para melhor desempenho.
Bloqueio de tabela durante a importação: A tabela de destino fica bloqueada durante toda a importação. Gravações DML e alterações DDL são bloqueadas; consultas de leitura não são afetadas.
Não destinado a gravações em tempo real: Para gravações streaming ou em tempo real, use o sink JDBC.
|
Parâmetro |
Descrição |
Obrigatório |
Tipo |
Padrão |
Observações |
|
|
Defina como |
Não |
STRING |
|
— |
|
|
Endereço IP ou nome do host do banco de dados OceanBase. |
Sim |
STRING |
— |
— |
|
|
Porta RPC do banco de dados OceanBase. |
Não |
INTEGER |
|
— |
|
|
Nome de usuário do banco de dados. |
Sim |
STRING |
— |
— |
|
|
Nome do tenant do OceanBase. |
Sim |
STRING |
— |
— |
|
|
Para tenants MySQL: nome do banco de dados. Para tenants Oracle: nome do proprietário. |
Sim |
STRING |
— |
— |
|
|
Nome da tabela de destino. |
Sim |
STRING |
— |
— |
|
|
Concorrência no lado do servidor para a tarefa de importação. O servidor limita o grau real de paralelismo com base nas especificações de CPU do tenant sem retornar erro. Fórmula: |
Não |
INTEGER |
|
— |
|
|
Número de registros armazenados em buffer antes de uma única gravação no OceanBase. |
Não |
INTEGER |
|
— |
|
|
Comportamento ao encontrar chaves primárias duplicadas. |
Não |
STRING |
|
— |
|
|
Modo de importação. |
Não |
STRING |
|
— |
|
|
Número máximo de linhas de erro toleradas. Linhas de erro incluem: chaves primárias duplicadas quando |
Não |
LONG |
|
— |
|
|
Timeout geral para a tarefa de importação bypass. |
Não |
DURATION |
|
— |
|
|
Timeout de heartbeat no lado do cliente. |
Não |
DURATION |
|
— |
|
|
Intervalo de heartbeat no lado do cliente. |
Não |
DURATION |
|
— |
Mapeamento de tipos
Modo compatível com MySQL
|
Tipo OceanBase |
Tipo Flink |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Modo compatível com Oracle
|
Tipo OceanBase |
Tipo Flink |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
Exemplos
Tabela de source e tabela de destino
O exemplo a seguir lê dados CDC de uma tabela de source OceanBase e os grava em uma tabela de destino JDBC. Uma definição de tabela de destino com importação bypass também está incluída para referência.
Todas as três tabelas usam 'connector' = 'oceanbase'. A source e o sink JDBC utilizam conjuntos de parâmetros diferentes; o sink direct-load define sink.mode = 'direct-load' e conecta-se via porta RPC.
-- OceanBase CDC source table (reads full history, then Binlog)
CREATE TEMPORARY TABLE oceanbase_source (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'oceanbase',
'hostname' = '<your-hostname>',
'port' = '3306',
'username' = '<your-username>',
'password' = '<your-password>',
'database-name' = '<your-database-name>',
'table-name' = '<your-table-name>'
);
-- OceanBase JDBC sink table (for real-time streaming writes)
CREATE TEMPORARY TABLE oceanbase_sink (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'oceanbase',
'url' = '<your-jdbc-url>',
'userName' = '<your-username>',
'password' = '<your-password>',
'tableName' = '<your-table-name>'
);
-- OceanBase bypass import sink table (for high-throughput batch writes)
-- Requires a bounded data source; set Flink to batch mode for best performance
CREATE TEMPORARY TABLE oceanbase_directload_sink (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'oceanbase',
'sink.mode' = 'direct-load',
'host' = '<your-host>',
'port' = '<your-rpc-port>',
'tenant-name' = '<your-tenant-name>',
'schema-name' = '<your-schema-name>',
'table-name' = '<your-table-name>',
'username' = '<your-username>',
'password' = '<your-password>'
);
BEGIN STATEMENT SET;
INSERT INTO oceanbase_sink
SELECT * FROM oceanbase_source;
END;
Tabela de dimensão
O exemplo a seguir faz uma junção entre uma source Datagen e uma tabela de dimensão OceanBase usando uma junção temporal. A política de cache ALL carrega toda a tabela de dimensão na memória antes do início do job.
CREATE TEMPORARY TABLE datagen_source (
a INT,
b BIGINT,
c STRING,
`proctime` AS PROCTIME()
) WITH (
'connector' = 'datagen'
);
-- OceanBase dimension table with ALL cache policy
-- ALL cache loads the full table at startup — suitable for small, stable tables
CREATE TEMPORARY TABLE oceanbase_dim (
a INT,
b VARCHAR,
c VARCHAR
) WITH (
'connector' = 'oceanbase',
'url' = '<your-jdbc-url>',
'userName' = '<your-username>',
'password' = '${secret_values.password}',
'tableName' = '<your-table-name>'
);
CREATE TEMPORARY TABLE blackhole_sink (
a INT,
b STRING
) WITH (
'connector' = 'blackhole'
);
INSERT INTO blackhole_sink
SELECT T.a, H.b
FROM datagen_source AS T
JOIN oceanbase_dim FOR SYSTEM_TIME AS OF T.`proctime` AS H
ON T.a = H.a;
Próximos passos
Supported connectors — lista completa de conectores disponíveis no Realtime Compute for Apache Flink
Visão geral do service Binlog do OceanBase — necessário para leituras incrementais CDC