Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Conector SQL do Hologres

Última atualização: Sep 07, 2026

Use o conector do Hologres para ler e gravar dados em tabelas do Hologres nos modos stream e batch, com suporte a CDC, consumo de binlog e atualizações parciais.

Visão geral

O Hologres é um data warehouse unificado em tempo real compatível com SQL padrão (PostgreSQL), OLAP em escala de petabytes, análises ad hoc e service de dados de alta concorrência e baixa latência. Ele se integra ao MaxCompute, Realtime Compute for Apache Flink e DataWorks. A tabela a seguir descreve as capacidades do conector.

Categoria

Descrição

Tipos suportados

Tabelas source, dimension e sink

Modo de execução

Modos stream e batch

Formato de dados

Não suportado

Métricas de monitoramento específicas

Métricas de monitoramento

  • Tabela source:

    • numRecordsIn

    • numRecordsInPerSecond

  • Tabela sink:

    • numRecordsOut

    • numRecordsOutPerSecond

    • currentSendTime

    Nota

    Detalhes das métricas: Monitoring metrics.

Tipos de API

DataStream e SQL

Suporte a atualizações ou exclusões em tabelas sink

Sim

Recursos

Recurso

Descrição

Real-time consumption of Hologres data

Permite ler dados do Hologres com ou sem binlog, tanto no modo CDC quanto no modo não-CDC.

Unified full and incremental consumption

Oferece suporte a consumo completo, incremental e unificado (completo e incremental).

Primary key conflict handling

Possibilita ignorar novos dados, substituir linhas inteiras ou atualizar apenas campos específicos.

Multi-stream merge and partial updates

Atualiza somente as colunas modificadas em vez de reescrever a linha inteira.

Consume binlog from partitioned tables (Beta)

Compatível com o consumo de binlog de tabelas particionadas físicas e lógicas. Em tabelas particionadas físicas, um único job pode monitorar todas as partições, incluindo as adicionadas recentemente.

Write to partitioned tables

Permite gravar em uma tabela pai para criar automaticamente as partições filhas correspondentes.

Sincronização em tempo real de tabela única ou banco de dados inteiro

Suporta sincronização em tempo real de uma única tabela ou de um banco de dados completo, com os seguintes recursos principais:

  • Detecção automática de evolução de schema da tabela source: Quando o schema de uma tabela source muda, o Hologres sincroniza essas alterações na tabela sink em tempo real.

  • Tratamento automático de alterações de schema: Ao ingerir novos dados, o Flink modifica primeiro o schema da tabela sink antes de gravar os dados.

Para mais informações, consulte Quick Start for real-time database synchronization.

Limitações e recomendações

Limitações

  • Tabelas externas não suportadas: O conector do Hologres não oferece suporte a Hologres foreign tables, como tabelas externas do MaxCompute.

  • Restrição de tipo de tempo: Atualmente, o conector do Hologres não suporta o consumo em tempo real de dados TIMESTAMP. Ao criar tabelas, use exclusivamente o tipo TIMESTAMPTZ .

  • Modo de varredura da tabela source (Ververica Runtime (VVR) 8 e anteriores): Por padrão, o conector lê dados no modo batch e realiza uma única varredura completa da tabela. Consequentemente, ele não consome dados adicionados posteriormente.

  • Restrição de watermark (Ververica Runtime (VVR) 8 e anteriores): O modo CDC não suporta definições de watermark. Para realizar agregações com janelas, use uma non-windowed aggregation solution.

  • Comportamento de funções de tempo no modo COPY: Ao usar o modo de gravação COPY_STREAM ou COPY_BULK_LOAD, os valores das colunas da tabela que usam CURRENT_TIMESTAMP ou NOW() como padrão são fixados no horário de início da conexão e não são atualizados para cada linha. A coluna de metadados de binlog hg_binlog_timestamp_us fornece o horário real de ingestão.

Recomendações

  • Seleção de formato de armazenamento:

    • Para consultas pontuais em tabelas dimension: Use armazenamento orientado a linhas. É necessário definir uma chave primária e uma chave de clustering.

    • Para consultas um-para-muitos em tabelas dimension: Use armazenamento orientado a colunas. Para obter desempenho ideal, configure adequadamente a distribution key e a segment key.

    • Para tabelas que exigem atualizações frequentes e consultas analíticas: Se uma tabela precisar suportar tanto o consumo de binlog em tempo real quanto análises OLAP, recomendamos fortemente o uso de armazenamento híbrido linha-coluna.

    Importante

    Caso você não especifique um formato de armazenamento ao criar uma tabela no Hologres, o sistema adotará o armazenamento orientado a colunas por padrão. Não é possível alterar o formato de armazenamento após a criação da tabela. Create a Table in Hologres | Table Storage Formats: Row-oriented, Column-oriented, and Row-column Hybrid.

  • Configurações de paralelismo do job: Defina o paralelismo do job Flink para corresponder ao número de shards na tabela do Hologres.

    -- In HoloWeb, run the following statement to find the number of shards in a table. Replace <tablename> with your table name.
    select tg.property_value from hologres.hg_table_properties tb join hologres.hg_table_group_properties tg on tb.property_value = tg.tablegroup_name where tb.property_key = 'table_group' and tg.property_key = 'shard_count' and table_name = '<tablename>';
  • Versões e recursos: Verifique regularmente o Hologres connector Release Notes para obter informações sobre problemas conhecidos, atualizações de recursos e compatibilidade de versões.

Notas de uso

  • Compatibilidade e limitações para modos de consumo do Hologres e VVR

    Tabela source

    • No VVR 8 e anteriores, selecione o modo de consumo usando o parâmetro sdkMode.

    • No VVR 11 e posteriores, selecione o modo de consumo usando o parâmetro source.binlog.read-mode.

    Versão do VVR

    Versão do Hologres

    Valor padrão/recomendado

    Modo de consumo real

    Observações

    ≥ 6.0.7

    < 2,0

    Personalizado

    holohub (padrão)

    Recomendamos definir o valor como jdbc.

    6.0.7 a 8.0.4

    ≥ 2,0

    jdbc (troca automática, sem necessidade de configuração)

    jdbc (forçado)

    O service holohub foi descontinuado no Hologres 2,0 e posteriores. O sistema troca automaticamente para o modo jdbc, o que pode causar problemas de permissão. Para configuração de permissões, consulte Permission issues.

    ≥ 8.0.5

    ≥ 2,1

    jdbc (troca automática, sem necessidade de configuração)

    jdbc (forçado)

    Sem problemas de permissão. No Hologres 2.1.27 e posteriores, o modo muda para jdbc_fixed.

    ≥ 11,1

    Qualquer versão

    AUTO (padrão)

    Selecionado automaticamente com base na versão do Hologres

    • Para Hologres 2.1.27 e posteriores, o modo jdbc é selecionado com conexões leves ativadas por padrão (o parâmetro connection.fixed.enabled é definido como true).

    • Para versões do Hologres de 2.1.0 a 2.1.26, o modo jdbc é selecionado.

    • Para Hologres 2,0 e anteriores, o modo holohub é selecionado.

    Importante

    No VVR 11,1 e posteriores, o conector consome dados de binlog por padrão. Certifique-se de ter enabled binlog. Caso contrário, podem ocorrer erros.

    Problemas de permissão

    Se você não for um superusuário, deverá conceder as permissões necessárias para consumir binlogs no modo JDBC.

    user_name é o ID da sua conta Alibaba Cloud ou usuário RAM (Account overview).

    -- In the standard PostgreSQL authorization model, grant the CREATE permission to the user and grant the replication role on the instance to the user.
    GRANT CREATE ON DATABASE <db_name> TO <user_name>;
    alter role <user_name> replication;
    
    -- If the database uses the simple permission model (SLMP), you cannot run the GRANT statement. Use spm_grant to grant the Admin permission on the database to the user. You can also grant the permission in the HoloWeb console.
    call spm_grant('<db_name>_admin', '<user_name>');
    alter role <user_name> replication;

    Tabela sink

    • No VVR 8 e anteriores, selecione o modo de gravação de dados usando o parâmetro sdkMode.

    • No VVR 11 e posteriores, selecione o modo de gravação de dados usando o parâmetro sink.write-mode.

    Versão do VVR

    Versão do Hologres

    Modo RPC afetado

    Modo de gravação real

    Valor padrão/recomendado

    Observações

    6.0.4 a 8.0.2

    < 2,0

    Não

    rpc

    Personalizado

    N/A

    6.0.4 a 8.0.2

    ≥ 2,0

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    Para evitar deduplicação, defina 'jdbcWriteBatchSize'='1'.

    ≥ 8.0.3

    Qualquer versão

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    Se você configurar o modo como rpc, o sistema mudará automaticamente para jdbc_fixed e definirá 'jdbcWriteBatchSize'='1' para evitar deduplicação.

    ≥ 8.0.5

    Qualquer versão

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    Se você configurar o modo como rpc, o sistema mudará automaticamente para jdbc_fixed e definirá 'deduplication.enabled'='false' para evitar deduplicação.

    Importante
    • O service rpc foi descontinuado no Hologres 2,0 e posteriores. Se você definir este parâmetro como rpc, o Flink alterará automaticamente o valor para jdbc_fixed. Se você definir o parâmetro com um valor diferente, o Flink usará o valor especificado.

    • O modo rpc foi removido no VVR 11,1 e posteriores. Recomendamos o uso do modo jdbc para conexões.

    • Para operações de gravação em cenários de alta concorrência, recomendamos o uso do modo jdbc_copy ou COPY_STREAM.

    Tabela dimension

    Versão do VVR

    Versão do Hologres

    Modo RPC afetado

    Modo de consumo real

    Valor padrão/recomendado

    Observações

    6.0.4 a 8.0.2

    < 2,0

    Não

    rpc

    Personalizado

    N/A

    6.0.4 a 8.0.2

    ≥ 2,0

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    O service rpc foi descontinuado para instâncias do Hologres versão 2,0 ou posterior. Se você definir este parâmetro como rpc, o Flink alterará automaticamente o valor para jdbc_fixed. No entanto, se você definir o parâmetro com um valor diferente, o Flink usará o valor especificado.

    ≥ 8.0.3

    Qualquer versão

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    ≥ 8.0.5

    Qualquer versão

    Sim

    jdbc_fixed (troca automática)

    Personalizado

    Importante

    O modo rpc foi removido no VVR 11,1 e posteriores. Por padrão, o modo jdbc é usado para conexões. Você também pode ativar o modo de conexão leve definindo o parâmetro connection.fixed.enabled.

  • Para ler dados JSONB de uma tabela source de binlog no modo JDBC, é necessário ativar um parâmetro GUC no nível do banco de dados.

    -- Enable the GUC parameter at the database level. Only a superuser can execute this command. You need to run this command only once for each database.
    alter database <db_name> set hg_experimental_enable_binlog_jsonb = on;
  • Uma operação UPDATE gera dois registros de binlog consecutivos: o registro update_before para os dados antigos, seguido pelo registro update_after para os novos dados.

  • Evite executar TRUNCATE ou outras operações de reconstrução de tabela em uma tabela source de binlog. Para mais informações, consulte FAQ.

  • Garanta que a precisão do tipo DECIMAL seja consistente entre o Flink e o Hologres para evitar erros. Para mais informações, consulte FAQ.

  • Ao usar o modo initial para consumo unificado de dados completos e incrementais de uma tabela source, a ordenação global não é garantida. Se os sistemas downstream dependerem de cálculos baseados em tempo, use um modo de consumo diferente, apenas de binlog.

Ativar binlog

Para novas tabelas

O recurso de leitura de dados em tempo real está desativado por padrão. Portanto, ao criar uma tabela no HoloWeb usando uma instrução DDL, você deve definir os parâmetros binlog.level e binlog.ttl. Veja um exemplo abaixo.

begin;
create table test_table(
  id int primary key, 
  title text not null, 
  body text);
call set_table_property('test_table', 'orientation', 'row');--Create a row-oriented table named test_table.
call set_table_property('test_table', 'clustering_key', 'id');--Create a clustering key on the id column.
call set_table_property('test_table', 'binlog.level', 'replica');--Enable the binlog feature.
call set_table_property('test_table', 'binlog.ttl', '86400');--Set the time-to-live (TTL) for the binlog in seconds.
commit;

Para tabelas existentes

No HoloWeb, use a seguinte instrução para ativar o Binlog em uma tabela existente e definir o TTL do Binlog. table_name é o nome da tabela para a qual você deseja ativar o Binlog.

-- Enable the binlog feature.
begin;
call set_table_property('<table_name>', 'binlog.level', 'replica');
commit;

-- Set the binlog time-to-live (TTL) in seconds.
begin;
call set_table_property('<table_name>', 'binlog.ttl', '2592000');
commit;

Parâmetros WITH

O VVR 11 renomeou ou removeu algumas opções do conector Hologres, mas mantém compatibilidade retroativa com o VVR 8. Consulte a documentação de parâmetros correspondente à sua versão.

Mapeamento de tipos

Para o mapeamento de tipos de dados entre Flink e Hologres, consulte Data type mapping between Flink and Hologres.

Nota

O Hologres suporta a sintaxe GENERATED ALWAYS AS para definir uma coluna gerada. Por exemplo:

ds TIMESTAMP NOT NULL GENERATED ALWAYS AS (date_trunc('month', create_time)) STORED

A restrição NOT NULL em uma coluna gerada é mapeada para um campo anulável, conforme esperado. Como o valor de uma coluna gerada é calculado pelo Hologres, o Flink não grava um valor para esse campo. Se a restrição NOT NULL fosse preservada, a operação de gravação falharia na validação do cliente Hologres. Esse comportamento não afeta a restrição NOT NULL em uma coluna comum (por exemplo, ds TIMESTAMP NOT NULL).

Exemplos

Exemplos de tabela source

Tabela source de binlog

Modo CDC

Neste modo, a source consome dados de binlog e define automaticamente o tipo RowKind apropriado do Flink para cada linha com base no hg_binlog_event_type, eliminando a necessidade de declarações explícitas. Esses tipos incluem INSERT, DELETE, UPDATE_BEFORE e UPDATE_AFTER. Isso permite o espelhamento dos dados da tabela, semelhante à funcionalidade Change Data Capture (CDC) em bancos de dados como MySQL ou PostgreSQL. Abaixo está um exemplo de DDL para a tabela source.

VVR 11+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='${secret_values.ak_id}',            --Use variables for your AccessKey pair to prevent key leakage. 
  'password'='${secret_values.ak_secret}',        
  'endpoint'='<yourEndpoint>',
  'source.binlog.change-log-mode'='ALL',  --Reads all changelog types, including INSERT, DELETE, UPDATE_BEFORE, and UPDATE_AFTER.
  'retry-count'='10',                     --The number of retries on a binlog read error.
  'retry-sleep-step-ms'='5000',           --The incremental backoff time between retries. The first retry waits 5 seconds, the second waits 10 seconds, and so on.
  'source.binlog.batch-size'='512'        --The batch size for reading binlog data.
);

VVR 8+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --Use variables for your AccessKey pair to prevent key leakage.   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'cdcMode' = 'true',
  'sdkMode'='jdbc',
  'binlogMaxRetryTimes' = '10',     --The number of retries on a binlog read error.
  'binlogRetryIntervalMs' = '500',  --The retry interval in milliseconds after a binlog read error.
  'binlogBatchReadSize' = '100'     --The batch size for reading binlog data.
);

Modo não-CDC

Neste modo, a source passa os dados de binlog consumidos para os nós downstream como dados regulares do Flink, tratando todos os registros como tipo INSERT. Você pode lidar com registros de um hg_binlog_event_type específico de acordo com seus requisitos de negócio. Abaixo está um exemplo de DDL para a tabela source.

VVR 11+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --Use variables for your AccessKey pair to prevent key leakage.   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'source.binlog.change-log-mode'='ALL_AS_APPEND_ONLY',  --All changelog types are treated as INSERT operations.
  'retry-count'='10',                     --The number of retries on a binlog read error.
  'retry-sleep-step-ms'='5000',           --The incremental backoff time between retries. The first retry waits 5 seconds, the second waits 10 seconds, and so on.
  'source.binlog.batch-size'='512'        --The batch size for reading binlog data.
);

VVR 8+

CREATE TEMPORARY TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --Use variables for your AccessKey pair to prevent key leakage.   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'binlogMaxRetryTimes' = '10',     --The number of retries on a binlog read error.
  'binlogRetryIntervalMs' = '500',  --The retry interval in milliseconds after a binlog read error.
  'binlogBatchReadSize' = '100'     --The batch size for reading binlog data.
);

Tabela source sem binlog

VVR 11+

Importante

A partir do VVR 11,1, o conector consome dados de binlog por padrão. Para ler de uma tabela source sem binlog, você deve definir explicitamente 'source.binlog' como 'false'. Para mais informações, consulte Binlog source table.

CREATE TEMPORARY TABLE hologres_source (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --Use variables for your AccessKey pair to prevent key leakage.   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'source.binlog'='false'                      --Specifies whether to consume binlog data.
);

VVR 8+

CREATE TEMPORARY TABLE hologres_source (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       --Use variables for your AccessKey pair to prevent key leakage.   
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'sdkMode' = 'jdbc'
);

Tabela sink

CREATE TEMPORARY TABLE datagen_source(
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='datagen'
);
CREATE TEMPORARY TABLE hologres_sink (
  name varchar, 
  age BIGINT,
  birthday BIGINT
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       -- Use variable management to prevent AK/SK key leakage.
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>'
);
INSERT INTO hologres_sink SELECT * from datagen_source;

Exemplo de tabela dimension

CREATE TEMPORARY TABLE datagen_source (
  a INT,
  b BIGINT,
  c STRING,
  proctime AS PROCTIME()
) WITH (
  'connector' = 'datagen'
);
CREATE TEMPORARY TABLE hologres_dim (
  a INT, 
  b VARCHAR, 
  c VARCHAR
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username' = '${secret_values.ak_id}',       -- Use variables for your AccessKey pair to prevent key leakage.
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>'
);
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 hologres_dim FOR SYSTEM_TIME AS OF T.proctime AS H ON T.a = H.a;

Recursos avançados

Ingestão unificada completa e incremental

Cenários

  • Este recurso aplica-se apenas a tabelas source que possuem chave primária. É recomendado para tabelas source do Hologres que utilizam o modo CDC.

  • O Hologres permite ativar o Binlog sob demanda. Você pode enable Binlog para tabelas existentes que já contêm dados.

Exemplo de código

VVR 11+

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'source.binlog.startup-mode' = 'INITIAL',   --Reads all historical data, then incrementally consumes the Binlog.
  'retry-count'='10',                         --Number of retries if an error occurs while reading Binlog data.
  'retry-sleep-step-ms'='5000',               --Incremental wait time between retries. The first retry waits 5 seconds, the second waits 10 seconds, and so on.
  'source.binlog.batch-size'='512'            --Number of rows to read from the Binlog in a single batch.
  );
Nota
  • Defina source.binlog.startup-mode como INITIAL para realizar uma leitura inicial completa da tabela antes de alternar para o consumo incremental do Binlog.

  • Se o parâmetro startTime estiver definido, seja diretamente ou selecionando um horário de início na interface de inicialização, ele terá precedência e definirá binlogStartUpMode como modo timestamp, substituindo outras configurações de modo, pois o parâmetro startTime tem prioridade mais alta.

VVR 8+

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn BIGINT,
  hg_binlog_event_type BIGINT,
  hg_binlog_timestamp_us BIGINT,
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='<yourTablename>',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'binlog' = 'true',
  'cdcMode' = 'true',
  'binlogStartUpMode' = 'initial', --Reads all historical data, then incrementally consumes the Binlog.
  'binlogMaxRetryTimes' = '10',     --Number of retries if an error occurs while reading Binlog data.
  'binlogRetryIntervalMs' = '500',  --Retry interval in milliseconds after a Binlog read error.
  'binlogBatchReadSize' = '100'     --Number of rows to read from the Binlog in a single batch.
  );
Nota
  • Defina binlogStartUpMode como initial para realizar uma leitura inicial completa da tabela antes de alternar para o consumo incremental do Binlog.

  • Se o parâmetro startTime estiver definido, seja diretamente ou selecionando um horário de início na interface de inicialização, ele terá precedência e definirá binlogStartUpMode como modo timestamp, substituindo outras configurações de modo, pois o parâmetro startTime tem prioridade mais alta.

Resolução de conflitos de chave primária

O conector oferece três estratégias para lidar com chaves primárias duplicadas durante a gravação.

VVR 11+

Especifique uma estratégia definindo o parâmetro sink.on-conflict-action.

Valor

Descrição

INSERT_OR_IGNORE

Mantém o primeiro registro recebido e ignora duplicatas subsequentes.

INSERT_OR_REPLACE

Sobrescreve o registro existente com o novo.

INSERT_OR_UPDATE (Padrão)

Atualiza apenas as colunas fornecidas na tabela sink, deixando as outras colunas do registro existente inalteradas.

VVR 8+

Especifique uma estratégia definindo o parâmetro mutatetype.

Valor

Descrição

insertorignore (Padrão)

Mantém o primeiro registro recebido e ignora duplicatas subsequentes.

insertorreplace

Sobrescreve o registro existente com o novo.

insertorupdate

Atualiza apenas as colunas fornecidas na tabela sink, deixando as outras colunas do registro existente inalteradas.

Por exemplo, suponha que uma tabela tenha as colunas a, b, c e d, sendo "a" a chave primária. Se a tabela sink fornecer apenas as colunas a e b, definir a estratégia como INSERT_OR_UPDATE atualizará apenas a coluna b. As colunas c e d permanecerão inalteradas.
Nota

No entanto, quaisquer colunas na tabela física que forem omitidas na tabela sink devem ser anuláveis. Caso contrário, a operação de gravação falhará.

Gravar em tabelas particionadas

Por padrão, o sink do Hologres grava dados em uma única tabela não particionada. Para gravar em uma tabela particionada direcionando para sua tabela pai, você deve ativar as seguintes opções.

VVR 11+

Para permitir que o conector crie automaticamente uma partição filha caso ela não exista, defina sink.create-missing-partition como true.

Nota
  • O VVR 11,1 e versões posteriores suportam gravação em tabelas particionadas por padrão e roteiam automaticamente os dados para as partições filhas corretas.

  • Defina o parâmetro tablename como o nome da tabela pai.

  • Se uma partição filha necessária não existir e sink.create-missing-partition=true não estiver definido, a operação de gravação falhará.

VVR 8+

  • Para rotear automaticamente os dados para as partições filhas correspondentes, defina partitionRouter como true.

  • Para permitir que o conector crie automaticamente uma partição filha caso ela não exista, defina createparttable como true.

Nota
  • Defina o parâmetro tablename como o nome da tabela pai.

  • Se uma partição filha necessária não existir e createparttable=true não estiver definido, a operação de gravação falhará.

Mesclar streams e atualizações parciais

Ao gravar múltiplos streams em uma única tabela wide do Hologres, o conector mescla registros com a mesma chave primária. As atualizações parciais gravam apenas as colunas modificadas em vez de substituir a linha inteira, melhorando o desempenho de gravação e a consistência dos dados.

Limitações

  • A tabela wide deve ter uma chave primária.

  • Cada stream de dados deve incluir todas as colunas que compõem a chave primária.

  • Para tabelas wide que usam armazenamento orientado a colunas, mesclar streams com uma alta taxa de requisições por segundo (RPS) pode causar alto uso de CPU. Para mitigar isso, considere desativar a codificação de dicionário para as colunas da tabela.

Exemplo

Suponha que você tenha dois streams de dados Flink. O primeiro stream contém as colunas a, b e c. O segundo stream contém as colunas a, d e e. A tabela wide do Hologres, WIDE_TABLE, contém as colunas a, b, c, d e e, onde a coluna a é a chave primária.

VVR 11+

// source1 and source2 are already defined.
CREATE TEMPORARY TABLE hologres_sink ( -- Declare columns a, b, c, d, and e.
  a BIGINT, 
  b STRING,
  c STRING,
  d STRING,
  e STRING,
  primary key(a) not enforced
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourWideTablename>',  -- The Hologres wide table, which contains columns a, b, c, d, and e.
  'username' = '${secret_values.ak_id}',
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'sink.on-conflict-action'='INSERT_OR_UPDATE',   -- Update specific columns based on the primary key.
  'sink.delete-strategy'='IGNORE_DELETE',         -- Strategy for handling retraction messages. IGNORE_DELETE is suitable for append-only or upsert streams where delete operations are not required.
  'sink.partial-insert.enabled'='true'            -- Enables partial updates. Only the columns specified in the INSERT statement are sent to the connector.
);

BEGIN STATEMENT SET;
INSERT INTO hologres_sink(a,b,c) select * from source1;  -- Declare that only columns a, b, and c are inserted.
INSERT INTO hologres_sink(a,d,e) select * from source2;  -- Declare that only columns a, d, and e are inserted.
END;

VVR 8+

// source1 and source2 are already defined.
CREATE TEMPORARY TABLE hologres_sink ( -- Declare columns a, b, c, d, and e.
  a BIGINT, 
  b STRING,
  c STRING,
  d STRING,
  e STRING,
  primary key(a) not enforced
) WITH (
  'connector'='hologres',           
  'dbname'='<yourDbname>',
  'tablename'='<yourWideTablename>',  -- The Hologres wide table, which contains columns a, b, c, d, and e.
  'username' = '${secret_values.ak_id}',
  'password' = '${secret_values.ak_secret}',
  'endpoint'='<yourEndpoint>',
  'mutatetype'='insertorupdate',    -- Update specific columns based on the primary key.
  'ignoredelete'='true',            -- Ignore DELETE requests generated by retraction messages.
  'partial-insert.enabled'='true'   -- Enable partial updates to update only the columns declared in the INSERT statement.
);

BEGIN STATEMENT SET;
INSERT INTO hologres_sink(a,b,c) select * from source1;  -- Declare that only columns a, b, and c are inserted.
INSERT INTO hologres_sink(a,d,e) select * from source2;  -- Declare that only columns a, d, and e are inserted.
END;
Nota

Defina ignoredelete como true para ignorar solicitações Delete geradas por mensagens de retração. Para VVR 8.0.8 e posteriores, recomendamos o uso de sink.delete-strategy para configurar várias estratégias para atualizações parciais.

Consumir binlog de tabelas particionadas (Beta)

O conector do Hologres suporta o consumo de binlog de tabelas particionadas físicas e lógicas. Para entender as diferenças entre elas, consulte CREATE LOGICAL PARTITION TABLE.

Consumir binlog de tabelas particionadas físicas

O conector do Hologres pode consumir binlog de uma tabela particionada e monitorar dinamicamente novas partições dentro de um único job. Isso melhora significativamente a eficiência e a usabilidade do processamento de dados em tempo real.

Notas de uso

  • Este recurso está disponível apenas para tabelas source de binlog no modo JDBC. Requer VVR 8.0.11 ou posterior e uma instância do Hologres versão 2.1.27 ou posterior.

  • O nome da partição deve seguir a convenção de nomenclatura dynamic partitioning: {parent_table}_{partition_value}. Partições fora do padrão podem não ser consumidas.

    Importante
    • Para o modo DYNAMIC, a versão 8.0.11 do VVR não suporta colunas de partição com delimitador - (como AAAA-MM-DD).

    • A partir do VVR 11,1, é possível consumir dados de colunas de partição que usam um formato personalizado.

    • Esta restrição não se aplica quando você grava em tabelas particionadas.

  • Ao declarar uma tabela source do Hologres no Flink, você deve incluir suas colunas de partição.

  • No modo DYNAMIC, uma tabela particionada deve ter o dynamic partitioning ativado. Além disso, o parâmetro de pré-criação de partição auto_partitioning.num_precreate deve ser maior que 1. Caso contrário, o job lançará uma exceção ao tentar consumir a partição mais recente.

  • No modo DYNAMIC, após a adição de uma nova partição, o conector deixa de consumir alterações de dados subsequentes de partições mais antigas.

Exemplos

Modo

Recursos

Descrição

DYNAMIC

Consumo dinâmico de partições

Monitora e consome automaticamente novas partições em ordem cronológica. Adequado para cenários de streaming de dados em tempo real.

STATIC

Consumo estático de partições

Consome apenas partições existentes (ou especificadas manualmente) e não descobre automaticamente novas partições. Adequado para processamento de dados históricos dentro de um intervalo fixo.

Modo dinâmico

VVR 11+

Suponha que uma tabela particionada do Hologres seja criada com o seguinte DDL, e que o binlog e o particionamento dinâmico estejam ativados.

CREATE TABLE "test_message_src1" (
    id int,
    title text,
    body text,
    dt text,
    PRIMARY KEY (id, dt)
)
PARTITION BY LIST (dt) WITH (
    binlog_level = 'replica', 
    auto_partitioning_enable =  'true',   -- Enables dynamic partitioning.
    auto_partitioning_time_unit = 'DAY',  -- Partitions are created daily. Example partition names: test_message_src1_20250512, test_message_src1_20250513.
    auto_partitioning_num_precreate = '2' -- Pre-creates two partitions.
);
-- For an existing partitioned table, you can also enable dynamic partitioning using ALTER TABLE.

No Flink, use a seguinte instrução SQL para consumir a tabela particionada test_message_src1 no modo DYNAMIC.

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  dt VARCHAR  -- The partition column of the Hologres partitioned table.
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src1',  -- The parent table for which dynamic partitioning is enabled.
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'source.binlog.partition-binlog-mode' = 'DYNAMIC', -- Dynamically monitors the latest partitions.
  'source.binlog.startup-mode' = 'initial'           -- Consumes all existing data, then consumes incremental data from the binlog.
);

VVR 8.0.11

Suponha que uma tabela particionada do Hologres seja criada com o seguinte DDL, e que o binlog e o particionamento dinâmico estejam ativados.

CREATE TABLE "test_message_src1" (
    id int,
    title text,
    body text,
    dt text,
    PRIMARY KEY (id, dt)
)
PARTITION BY LIST (dt) WITH (
    binlog_level = 'replica', 
    auto_partitioning_enable =  'true',   -- Enables dynamic partitioning.
    auto_partitioning_time_unit = 'DAY',  -- Partitions are created daily. Example partition names: test_message_src1_20241027, test_message_src1_20241028.
    auto_partitioning_num_precreate = '2' -- Pre-creates two partitions.
);

-- For an existing partitioned table, you can also enable dynamic partitioning using ALTER TABLE.

No Flink, use a seguinte instrução SQL para consumir dados da tabela particionada test_message_src1 no modo DYNAMIC.

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  dt VARCHAR  -- The partition column of the Hologres partitioned table.
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src1',  -- The parent table for which dynamic partitioning is enabled.
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'binlog' = 'true',
  'partition-binlog.mode' = 'DYNAMIC',  -- Dynamically monitors the latest partitions.
  'binlogstartUpMode' = 'initial',      -- Consumes all existing data, then consumes incremental data from the binlog.
  'sdkMode' = 'jdbc_fixed'              -- Use this mode to avoid connection limit issues.
);

Modo estático

VVR 11+

Suponha que uma tabela particionada do Hologres seja criada com o seguinte DDL, e que o binlog esteja ativado.

CREATE TABLE test_message_src2 (
    id int,
    title text,
    body text,
    color text,
    PRIMARY KEY (id, color)
)
PARTITION BY LIST (color) WITH (
    binlog_level = 'replica'
);
create table test_message_src2_red partition of test_message_src2 for values in ('red');
create table test_message_src2_blue partition of test_message_src2 for values in ('blue');
create table test_message_src2_green partition of test_message_src2 for values in ('green');
create table test_message_src2_black partition of test_message_src2 for values in ('black');

No Flink, use a seguinte instrução SQL para consumir a tabela particionada test_message_src2 no modo STATIC.

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  color VARCHAR  -- The partition column of the Hologres partitioned table.
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src2',  -- The partitioned table.
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'source.binlog.partition-binlog-mode' = 'STATIC', -- Consumes a fixed set of partitions.
  'source.binlog.partition-values-to-read' = 'red,blue,green',  -- Consumes only the three specified partitions. The 'black' partition is not consumed. New partitions are also not consumed. If this option is not set, the job consumes all partitions of the parent table.
  'source.binlog.startup-mode' = 'initial'  -- Consumes all existing data, then consumes incremental data from the binlog.
);

VVR 8.0.11

Suponha que uma tabela particionada do Hologres seja criada com o seguinte DDL, e que o binlog esteja ativado.

CREATE TABLE test_message_src2 (
    id int,
    title text,
    body text,
    color text,
    PRIMARY KEY (id, color)
)
PARTITION BY LIST (color) WITH (
    binlog_level = 'replica'
);
create table test_message_src2_red partition of test_message_src2 for values in ('red');
create table test_message_src2_blue partition of test_message_src2 for values in ('blue');
create table test_message_src2_green partition of test_message_src2 for values in ('green');
create table test_message_src2_black partition of test_message_src2 for values in ('black');

No Flink, use a seguinte instrução SQL para consumir dados da tabela particionada test_message_src2 no modo STATIC.

CREATE TEMPORARY TABLE hologres_source
(
  id INTEGER,
  title VARCHAR,
  body VARCHAR,
  color VARCHAR  -- The partition column of the Hologres partitioned table.
)
with (
  'connector' = 'hologres',
  'dbname' = '<yourDatabase>',
  'tablename' = 'test_message_src2',  -- The partitioned table.
  'username' = '<yourUserName>',
  'password' = '<yourPassword>',
  'endpoint' = '<yourEndpoint>',
  'binlog' = 'true',
  'partition-binlog.mode' = 'STATIC', -- Consumes a fixed set of partitions.
  'partition-values-to-read' = 'red,blue,green',  -- Consumes only the three specified partitions. The 'black' partition is not consumed. New partitions are also not consumed. If this option is not set, the job consumes all partitions of the parent table.
  'binlogstartUpMode' = 'initial',  -- Consumes all existing data, then consumes incremental data from the binlog.
  'sdkMode' = 'jdbc_fixed' -- Use this mode to avoid connection limit issues.
);

Consumir binlog de tabelas particionadas lógicas

O conector do Hologres suporta o consumo de binlog de tabelas particionadas lógicas e permite especificar quais partições consumir por meio de opções.

Notas de uso

  • O consumo de binlog de partições específicas de uma tabela particionada lógica requer VVR 11.0.0 ou posterior e uma instância do Hologres versão V3.1 ou posterior.

  • O consumo de binlog de todas as partições de uma tabela particionada lógica segue a mesma abordagem de uma tabela não particionada (Source tables).

Exemplos

Parâmetro

Descrição

Exemplo

source.binlog.logical-partition-filter-column-names

Os nomes das colunas de partição que especificam quais partições consumir. Coloque os nomes das colunas entre aspas duplas ("). Separe vários nomes de coluna com vírgulas (,). Se um nome de coluna contiver aspas duplas, escape-as com aspas duplas adicionais.

'source.binlog.logical-partition-filter-column-names'='"Pt","id"'

Duas colunas de partição são usadas: Pt e id.

source.binlog.logical-partition-filter-column-values

Os valores de partição que especificam quais partições consumir. Uma partição é especificada por um conjunto de valores, um para cada coluna de partição. Coloque cada valor entre aspas duplas ("). Separe os valores da mesma partição com vírgula (,). Separe as partições com ponto e vírgula (;). Se um valor contiver aspas duplas, escape-as com aspas duplas adicionais.

'source.binlog.logical-partition-filter-column-values'='"20240910","0";"special""value","9"'

Isso especifica duas partições para consumo. O valor da primeira partição é (20240910, 0) e o da segunda é (special"value, 9).

Suponha que você tenha criado a seguinte tabela no Hologres.

CREATE TABLE holo_table (
    id int not null,
    name text,
    age numeric(18,4),
    "Pt" text,
    primary key(id, "Pt")
)
LOGICAL PARTITION BY LIST ("Pt", id)
WITH (
    binlog_level ='replica'
);

Para consumir o binlog desta tabela no Flink:

CREATE TEMPORARY TABLE test_src_binlog_table(
  id INTEGER,
  name VARCHAR,
  age decimal(18,4),
  `Pt` VARCHAR
) WITH (
  'connector'='hologres',
  'dbname'='<yourDbname>',
  'tablename'='holo_table',
  'username'='<yourAccessID>',
  'password'='<yourAccessSecret>',
  'endpoint'='<yourEndpoint>',
  'source.binlog'='true',
  'source.binlog.logical-partition-filter-column-names'='"Pt","id"',
  'source.binlog.logical-partition-filter-column-values'='<yourPartitionColumnValues>',
  'source.binlog.change-log-mode'='ALL',  --Reads all changelog types, including INSERT, DELETE, UPDATE_BEFORE, and UPDATE_AFTER.
  'retry-count'='10',                     -- The number of retries for binlog read errors.
  'retry-sleep-step-ms'='5000',           --The incremental backoff time between retries. The first retry waits 5 seconds, the second waits 10 seconds, and so on.
  'source.binlog.batch-size'='512'        --The number of rows to read from the binlog in a single batch.
);

API DataStream

Importante

Para ler ou gravar no Hologres usando a API DataStream, configure o conector DataStream correspondente (How to use DataStream connectors). O conector DataStream do Hologres está disponível no Maven Central. Para depuração local, use o Uber JAR (Run and debug jobs that contain connectors locally).

Tabela source do Hologres

Tabela source de binlog

O VVR fornece a classe HologresBinlogSource para ler dados de binlog do Hologres. O exemplo a seguir mostra como construir uma source de binlog do Hologres.

VVR 11.3+

Importante

A partir do VVR 11.1.2, os parâmetros JDBCOptions e startTimeMs foram removidos do construtor HologresBinlogSource. A partir do VVR 11.3, um parâmetro List<Subscribe.BinlogFilter> foi adicionado. Se você usar o VVR 11 ou posterior, recomendamos o uso do VVR 11.3 ou superior.

public class Sample {                                                                                                                                                                          
        public static void main(String[] args) throws Exception {
            final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            // Initialize the schema for the table to read. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
            TableSchema schema = TableSchema.builder()
                    .field("a", DataTypes.INT())
                    .field("b", DataTypes.STRING())
                    .field("c", DataTypes.TIMESTAMP())
                    .build();

            // The name of the table to read.
            String sourceTableName = "sourceTableName";

            // Parameters for the Hologres connection.
            Configuration config = new Configuration();
            config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
            config.setString(HologresConfigs.USERNAME, "yourUserName");
            config.setString(HologresConfigs.PASSWORD, "yourPassword");
            config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
            config.setString(HologresConfigs.TABLE, sourceTableName);
            config.set(HologresConfigs.BINLOG, true);
            config.set(HologresConfigs.BINLOG_CHANGE_LOG_MODE, BinlogChangeLogMode.ALL);
            // Build the Hologres binlog source.
            HologresBinlogSource source = new HologresBinlogSource(
                    new HologresConnectionParam(config),
                    schema,
                    config,
                    StartupMode.INITIAL,
                    sourceTableName,
                    "",
                    Collections.emptyList(),
                    -1,
                    Collections.emptySet(),
                    Collections.emptyList()
            );
            env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
            env.execute();
        }
  }

VVR 8.0.11+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Initialize the schema for the table to read. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Parameters for the Hologres connection.
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // Build JDBCOptions.
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // Build the Hologres binlog source.
        long startTimeMs = 0;
        HologresBinlogSource source = new HologresBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.INITIAL,
                "",
                "",
                -1,
                Collections.emptySet(),
                new ArrayList<>()
        );
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}

VVR 8.0.7+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Initialize the schema for the table to read. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Parameters for the Hologres connection.
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // Build JDBCOptions.
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // Build the Hologres binlog source.
        long startTimeMs = 0;
        HologresBinlogSource source = new HologresBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.INITIAL,
                "",
                "",
                -1,
                Collections.emptySet()
        );
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}

VVR 6.0.7+

public class Sample {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Initialize the schema for the table to read. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .build();
         // Parameters for the Hologres connection.
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        config.setBoolean(HologresBinlogConfigs.OPTIONAL_BINLOG, true);
        config.setBoolean(HologresBinlogConfigs.BINLOG_CDC_MODE, true);
        // Build JDBCOptions.
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        // Set or create the default slot name.
        config.setString(HologresBinlogConfigs.JDBC_BINLOG_SLOT_NAME, HoloBinlogUtil.getOrCreateDefaultSlotForJDBCBinlog(jdbcOptions));

        boolean cdcMode = config.get(HologresBinlogConfigs.BINLOG_CDC_MODE) && config.get(HologresBinlogConfigs.OPTIONAL_BINLOG);
        // Build the JDBCBinlogRecordConverter.
        JDBCBinlogRecordConverter recordConverter = new JDBCBinlogRecordConverter(
                jdbcOptions.getTable(),
                schema,
                new HologresConnectionParam(config),
                cdcMode,
                Collections.emptySet());
        
        // Build the Hologres binlog source.
        long startTimeMs = 0;
        HologresJDBCBinlogSource source = new HologresJDBCBinlogSource(
                new HologresConnectionParam(config),
                schema,
                config,
                jdbcOptions,
                startTimeMs,
                StartupMode.TIMESTAMP,
                recordConverter,
                "",
                -1);
        env.fromSource(source, WatermarkStrategy.noWatermarks(), "Test source").print();
        env.execute();
    }
}
Importante

Se você usar uma versão do motor Flink anterior a 8.0.5 ou Hologres anterior a 2,1, certifique-se de que o usuário seja um superusuário ou tenha a Função de Replicação (Hologres permission issues).

Tabela source sem binlog

O VVR fornece a classe HologresBulkreadInputFormat, uma implementação de RichInputFormat, para ler dados de tabelas do Hologres. O exemplo a seguir mostra como construir uma source do Hologres.

public class Sample {
    public static void main(String[] args) throws Exception {
        // set up the Java DataStream API
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Initialize the schema for the table to read. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .field("c", DataTypes.TIMESTAMP())
                .build();
        // Parameters for the Hologres connection.
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        // Build JDBCOptions.
        JDBCOptions jdbcOptions = JDBCUtils.getJDBCOptions(config);
        HologresBulkreadInputFormat inputFormat = new HologresBulkreadInputFormat(
                new HologresConnectionParam(config),
                jdbcOptions,
                schema,
                "",
                -1);
        TypeInformation<RowData> typeInfo = InternalTypeInfo.of(schema.toRowDataType().getLogicalType());
        env.addSource(new InputFormatSourceFunction<>(inputFormat, typeInfo)).returns(typeInfo).print();
        env.execute();
    }
}

Dependência Maven

O conector DataStream do Hologres está disponível no Maven Central.

<dependency>
    <groupId>com.alibaba.ververica</groupId>
    <artifactId>ververica-connector-hologres</artifactId>
    <version>${vvr-version}</version>
</dependency>

Tabela sink do Hologres

VVR 11+

public class Sample {
      public static void main(String[] args) throws Exception {
          final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
          // Initialize the schema for the table to write to. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
          TableSchema tableSchema = TableSchema.builder()
                  .field("a", DataTypes.INT().notNull())
                  .field("b", DataTypes.STRING())
                  .primaryKey("a")
                  .build();
          // Parameters for the Hologres connection.
          Configuration config = new Configuration();
          config.set(HologresConfigs.ENDPOINT, "yourEndpoint");
          config.set(HologresConfigs.USERNAME, "yourUserName");
          config.set(HologresConfigs.PASSWORD, "yourPassword");
          config.set(HologresConfigs.DATABASE, "yourDatabaseName");
          config.set(HologresConfigs.TABLE, "yourTableName");
          HologresConnectionParam connectionParam = new HologresConnectionParam(config);
          HologresTableSchema hologresTableSchema =
                  HologresTableSchema.get(connectionParam.getJDBCOptions());
          // The indexes of the columns to write to the sink.
          Integer[] targetColumnIndexes = {0, 1};
          // Build the Hologres sink.
          HologresSinkFunction sinkFunction =
                  new HologresSinkFunction(
                          connectionParam, tableSchema, targetColumnIndexes, hologresTableSchema);
          TypeInformation<RowData> typeInfo = InternalTypeInfo.of(tableSchema.toRowDataType().getLogicalType());
          env.fromElements((RowData) GenericRowData.of(101, StringData.fromString("name"))).returns(typeInfo).addSink(sinkFunction);
          env.execute();
      }
  }

VVR 8+

public class Sample {
    public static void main(String[] args) throws Exception {
        // set up the Java DataStream API
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // Initialize the schema for the table to write to. The schema must match the fields of the Hologres table schema. You can define a subset of the fields.
        TableSchema schema = TableSchema.builder()
                .field("a", DataTypes.INT())
                .field("b", DataTypes.STRING())
                .build();
        // Parameters for the Hologres connection.
        Configuration config = new Configuration();
        config.setString(HologresConfigs.ENDPOINT, "yourEndpoint");
        config.setString(HologresConfigs.USERNAME, "yourUserName");
        config.setString(HologresConfigs.PASSWORD, "yourPassword");
        config.setString(HologresConfigs.DATABASE, "yourDatabaseName");
        config.setString(HologresConfigs.TABLE, "yourTableName");
        config.setString(HologresConfigs.SDK_MODE, "jdbc");
        HologresConnectionParam hologresConnectionParam = new HologresConnectionParam(config);
        
         // Build a Hologres writer to write data as RowData.
        AbstractHologresWriter<RowData> hologresWriter = HologresJDBCWriter.createRowDataWriter(
                hologresConnectionParam, 
                schema, 
                HologresTableSchema.get(hologresConnectionParam), 
                new Integer[0]);
        // Build the Hologres sink.
        HologresSinkFunction sinkFunction = new HologresSinkFunction(hologresConnectionParam, hologresWriter);
        TypeInformation<RowData> typeInfo = InternalTypeInfo.of(schema.toRowDataType().getLogicalType());
        env.fromElements((RowData) GenericRowData.of(101, StringData.fromString("name"))).returns(typeInfo).addSink(sinkFunction);
        env.execute();
    }
}

Colunas de metadados

O VVR 8.0.11 e posteriores suportam colunas de metadados em tabelas source de binlog. Declare campos de binlog como hg_binlog_event_type como colunas de metadados para acessar o nome do banco de dados source, nome da tabela, tipo de alteração e timestamp do evento para lógica personalizada (por exemplo, filtragem de eventos DELETE).

Parâmetro

Tipo

Descrição

db_name

STRING NOT NULL

O nome do banco de dados que contém a linha.

table_name

STRING NOT NULL

O nome da tabela que contém a linha.

hg_binlog_lsn

BIGINT NOT NULL

Um campo de sistema para o número de sequência do binlog. O valor aumenta monotonicamente, mas não contiguamente dentro de um shard. A unicidade e a ordem não são garantidas entre shards.

hg_binlog_timestamp_us

BIGINT NOT NULL

O timestamp do evento de alteração no banco de dados, em microssegundos (us).

hg_binlog_event_type

BIGINT NOT NULL

O tipo de alteração da linha. Os valores válidos são:

  • 5: Uma mensagem INSERT.

  • 2: Uma mensagem DELETE.

  • 3: A imagem da linha antes de uma operação UPDATE.

  • 7: A imagem da linha após uma operação UPDATE.

hg_shard_id

INT NOT NULL

O ID do shard de dados que contém a linha (table group and shard).

Em uma instrução DDL, você pode declarar uma coluna de metadados usando <meta_column_name> <datatype> METADATA VIRTUAL. Veja um exemplo:

CREATE TABLE test_message_src_binlog_table(
  hg_binlog_lsn bigint METADATA VIRTUAL
  hg_binlog_event_type bigint METADATA VIRTUAL
  hg_binlog_timestamp_us bigint METADATA VIRTUAL
  hg_shard_id int METADATA VIRTUAL
  db_name string METADATA VIRTUAL
  table_name string METADATA VIRTUAL
  id INTEGER,
  title VARCHAR,
  body VARCHAR
) WITH (
  'connector'='hologres',
  ...
  );

Perguntas frequentes

Referências