Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:OceanBase

Última atualização: Sep 03, 2026

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:

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 INTO

  • Com 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

connector

Defina como oceanbase.

Sim

STRING

password

Senha do banco de dados.

Sim

STRING

Parâmetros da tabela de source

Importante

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 earliest-offset, specific-offset e timestamp, o esquema da tabela não deve mudar entre a posição especificada do Binlog e o início do job.

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

URL JDBC. Deve incluir o nome do banco de dados MySQL ou o nome do service Oracle.

Sim

STRING

userName

Nome de usuário do banco de dados.

Sim

STRING

cache

Política de cache para consultas à tabela de dimensão.

Não

STRING

ALL

ALL: Carrega todos os dados antes do início do job; recarrega após expiração. Adequado para tabelas pequenas com muitas falhas de consulta. Aumente a memória do nó de junção para pelo menos o dobro do tamanho da tabela para suportar carregamento assíncrono. LRU: Armazena em cache um subconjunto de linhas; requer cacheSize. None: Sem cache.

cacheSize

Número máximo de entradas em cache.

Não

INTEGER

100000

Obrigatório quando cache=LRU. Ignorado quando cache=ALL.

cacheTTLMs

Timeout do cache em milissegundos. O comportamento depende da configuração de cache: para LRU, as entradas expiram após esta duração (sem expiração por padrão); para ALL, o cache completo é recarregado após esta duração (sem recarga por padrão); para None, este parâmetro não tem efeito.

Não

LONG

Long.MAX_VALUE

maxRetryTimeout

Duração máxima de nova tentativa.

Não

DURATION

60s

Parâmetros da tabela de destino (JDBC)

Parâmetro

Descrição

Obrigatório

Tipo

Padrão

Observações

url

URL JDBC. Deve incluir o nome do banco de dados MySQL ou o nome do service Oracle.

Sim

STRING

userName

Nome de usuário do banco de dados.

Sim

STRING

tableName

Nome da tabela de destino.

Sim

STRING

sink.mode

Modo de gravação. Defina como jdbc para gravações padrão; defina como direct-load para importação bypass.

Sim

STRING

jdbc

compatibleMode

Modo de compatibilidade do OceanBase. Valores válidos: mysql, oracle.

Não

STRING

mysql

Parâmetro específico do OceanBase.

maxRetryTimes

Número máximo de novas tentativas de gravação.

Não

INTEGER

3

poolInitialSize

Tamanho inicial do pool de conexões.

Não

INTEGER

1

poolMaxActive

Número máximo de conexões ativas no pool.

Não

INTEGER

8

poolMaxWait

Tempo máximo de espera (ms) por uma conexão do pool.

Não

INTEGER

2000

poolMinIdle

Número mínimo de conexões ociosas no pool.

Não

INTEGER

1

connectionProperties

Propriedades de conexão JDBC no formato k1=v1;k2=v2.

Não

STRING

ignoreDelete

Define se operações de exclusão devem ser ignoradas.

Não

BOOLEAN

false

excludeUpdateColumns

Colunas a excluir das atualizações, separadas por vírgula (por exemplo, column1,column2). As colunas de chave primária são sempre excluídas independentemente desta configuração.

Não

STRING

partitionKey

Chave de partição. Quando definida, os dados são agrupados por esta chave antes da aplicação da modRule.

Não

STRING

modRule

Regra de agrupamento no formato column_name mod number (por exemplo, user_id mod 8). A coluna deve ser numérica. Os dados são primeiro particionados por partitionKey e depois agrupados dentro de cada partição por esta regra.

Não

STRING

bufferSize

Tamanho do buffer de dados (número de registros).

Não

INTEGER

1000

flushIntervalMs

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

1000

retryIntervalMs

Intervalo de nova tentativa (ms).

Não

INTEGER

5000

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

sink.mode

Defina como direct-load para usar importação bypass.

Não

STRING

jdbc

host

Endereço IP ou nome do host do banco de dados OceanBase.

Sim

STRING

port

Porta RPC do banco de dados OceanBase.

Não

INTEGER

2882

username

Nome de usuário do banco de dados.

Sim

STRING

tenant-name

Nome do tenant do OceanBase.

Sim

STRING

schema-name

Para tenants MySQL: nome do banco de dados. Para tenants Oracle: nome do proprietário.

Sim

STRING

table-name

Nome da tabela de destino.

Sim

STRING

parallel

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: MIN(tenant_cores × 2, parallel) × partition_nodes. Por exemplo, com 2 núcleos de CPU, parallel=10 e 2 nós de partição: MIN(4, 10) × 2 = 8.

Não

INTEGER

8

buffer-size

Número de registros armazenados em buffer antes de uma única gravação no OceanBase.

Não

INTEGER

1024

dup-action

Comportamento ao encontrar chaves primárias duplicadas. STOP_ON_DUP: falha na importação. REPLACE: sobrescreve a linha existente. IGNORE: descarta a linha recebida.

Não

STRING

REPLACE

load-method

Modo de importação. full: importação bypass padrão. inc: modo incremental, verifica conflitos de chave primária (observer 4.3.2+, dup-action=REPLACE não suportado). inc_replace: modo de substituição incremental, sobrescreve linhas existentes diretamente sem verificações de conflito (observer 4.3.2+, dup-action é ignorado).

Não

STRING

full

max-error-rows

Número máximo de linhas de erro toleradas. Linhas de erro incluem: chaves primárias duplicadas quando dup-action=STOP_ON_DUP, contagens de colunas incompatíveis e linhas onde a conversão de tipo falha.

Não

LONG

0

timeout

Timeout geral para a tarefa de importação bypass.

Não

DURATION

7d

heartbeat-timeout

Timeout de heartbeat no lado do cliente.

Não

DURATION

60s

heartbeat-interval

Intervalo de heartbeat no lado do cliente.

Não

DURATION

10s

Mapeamento de tipos

Modo compatível com MySQL

Tipo OceanBase

Tipo Flink

TINYINT

TINYINT

SMALLINT, TINYINT UNSIGNED

SMALLINT

INT, MEDIUMINT, SMALLINT UNSIGNED

INT

BIGINT, INT UNSIGNED

BIGINT

BIGINT UNSIGNED

DECIMAL(20, 0)

REAL, FLOAT

FLOAT

DOUBLE

DOUBLE

NUMERIC(p, s), DECIMAL(p, s)

DECIMAL(p, s) (p ≤ 38)

BOOLEAN, TINYINT(1)

BOOLEAN

DATE

DATE

TIME [(p)]

TIME [(p)] [WITHOUT TIME ZONE]

DATETIME [(p)], TIMESTAMP [(p)]

TIMESTAMP [(p)] [WITHOUT TIME ZONE]

CHAR(n)

CHAR(n)

VARCHAR(n)

VARCHAR(n)

BIT(n)

BINARY(⌈n/8⌉)

BINARY(n)

BINARY(n)

VARBINARY(N)

VARBINARY(N)

TINYTEXT, TEXT, MEDIUMTEXT, LONGTEXT

STRING

TINYBLOB, BLOB, MEDIUMBLOB, LONGBLOB

BYTES (máx. 2.147.483.647 bytes)

Modo compatível com Oracle

Tipo OceanBase

Tipo Flink

NUMBER(p, s≤0), p−s < 3

TINYINT

NUMBER(p, s≤0), p−s < 5

SMALLINT

NUMBER(p, s≤0), p−s < 10

INT

NUMBER(p, s≤0), p−s < 19

BIGINT

NUMBER(p, s≤0), 19 ≤ p−s ≤ 38

DECIMAL(p−s, 0)

NUMBER(p, s>0)

DECIMAL(p, s)

NUMBER(p, s≤0), p−s > 38

STRING

FLOAT, BINARY_FLOAT

FLOAT

BINARY_DOUBLE

DOUBLE

NUMBER(1)

BOOLEAN

DATE, TIMESTAMP [(p)]

TIMESTAMP [(p)] [WITHOUT TIME ZONE]

CHAR(n), NCHAR(n), NVARCHAR2(n), VARCHAR(n), VARCHAR2(n), CLOB

STRING

BLOB, ROWID

BYTES

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