Todos os produtos
Search
Central de documentação

Hologres:Consumir Binlog via JDBC

Última atualização: Jun 28, 2026

O Binlog do Hologres rastreia operações INSERT, UPDATE, DELETE e TRUNCATE como um fluxo de eventos de alteração. Esta página demonstra como consumir esse fluxo usando Java Database Connectivity (JDBC) diretamente ou por meio do Holo-Client, o SDK Java do Hologres.

Pré-requisitos

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

  • Binlog ativado na tabela de destino. Consulte Assinar dados do Binlog do Hologres.

  • A extensão hg_binlog criada (apenas para Hologres V1.x). A partir da V2.0, a extensão é nativa e não requer configuração manual. Para versões do Hologres anteriores à V2.0, um superusuário deve executar a seguinte instrução uma vez por banco de dados. Se criar um novo banco de dados, execute-a novamente.

    Importante

    Nunca use DROP EXTENSION <extension_name> CASCADE. A opção CASCADE remove todos os dados da extensão — incluindo PostGIS, RoaringBitmap, Proxima, Binlog e BSI — e exclui todos os objetos dependentes, como tabelas, visualizações, metadados e dados do servidor.

    -- Create the extension.
    CREATE EXTENSION hg_binlog;
    
    -- Remove the extension.
    DROP EXTENSION hg_binlog;
  • Acesso configurado conforme sua versão do Hologres: para criar uma publicação e um slot de replicação (necessário para todas as versões que utilizam este método), o usuário deve ter um dos seguintes conjuntos de permissões:

    • Permissões de superusuário na instância

    • Permissões de proprietário na tabela de destino, permissões CREATE DATABASE e a Replication Role na instância

    Nota

    A partir da V2.0, o modo Holohub não é mais suportado. Antes de atualizar para a V2.0 ou superior, atualize o Flink para a versão 8.0.5. Após a atualização, o modo JDBC será utilizado automaticamente.

    Versão do Hologres

    Versão do mecanismo Flink

    Requisito

    V2.1 e posteriores

    8.0.5 e posteriores

    Permissões de leitura na tabela de destino. Nenhum slot de replicação necessário.

    V2.0

    8.0.5 e anteriores

    Uma publicação e um slot de replicação. Consulte Criar uma publicação e um slot de replicação.

    V1.3 e anteriores

    8.0.5 e anteriores

    Permissões de leitura na tabela de destino. O modo Holohub é usado automaticamente.

Limitações

  • O consumo de Binlog baseado em JDBC requer Hologres V1.1 ou posterior.

  • Tipos de dados de coluna suportados: INTEGER, BIGINT, SMALLINT, TEXT, CHAR(n), VARCHAR(n), REAL, DOUBLE PRECISION, BOOLEAN, NUMERIC(38,8), DATE, TIME, TIMETZ, TIMESTAMP, TIMESTAMPTZ, BYTEA, JSON, SERIAL, OID, int4[], int8[], float4[], float8[], boolean[] e text[]. Se uma tabela contiver colunas de tipos não suportados, o consumo falhará. A partir da V1.3.36, colunas JSONB também são suportadas. Ative o seguinte parâmetro Grand Unified Configuration (GUC) antes de iniciar o consumo:

    -- Enable at the session level.
    SET hg_experimental_enable_binlog_jsonb = ON;
    
    -- Enable at the database level.
    ALTER DATABASE <db_name> SET hg_experimental_enable_binlog_jsonb = ON;
  • Cada shard de cada tabela consumida utiliza uma conexão Walsender. As conexões Walsender são independentes das conexões regulares.

  • O número máximo de Walsenders por nó frontend varia conforme a versão. Para verificar o limite atual, execute:

    1. Identifique jobs ativos de consumo de Binlog via JDBC e interrompa os desnecessários.

    2. Verifique se as configurações de grupo de tabelas e contagem de shards estão adequadas. Consulte Melhores práticas para configurar grupos de tabelas.

    3. Se o limite ainda for excedido, escale horizontalmente a instância.

    Versão

    Máximo de Walsenders por nó frontend

    V2.2 e posteriores

    600

    V2.0 e V2.1

    1.000

    V1.1.26 até V2.0

    100

    SHOW max_wal_senders;

    Valores padrão: o limite total é max_wal_senders × número de nós frontend. Para consultar o número de nós frontend de uma determinada especificação de instância, veja Gerenciamento de instâncias. Use esta fórmula para estimar quantas tabelas podem ser consumidas simultaneamente:

    Number of tables ≤ (max_wal_senders × number of frontend nodes) / table shard count

    Exemplo: Uma instância com dois nós frontend executando a V2.2 tem um total de 600 × 2 = 1.200 Walsenders. Para tabelas com contagem de shards igual a 20, ela suporta até 60 consumos simultâneos de tabelas. Se dois jobs consumirem da mesma tabela ao mesmo tempo, cada um utilizará 20 Walsenders — totalizando 40 no limite geral. Ao atingir o limite de Walsenders, a seguinte mensagem será exibida:

    FATAL: sorry, too many wal senders already

    Para resolver isso:

  • Instâncias secundárias somente leitura: o consumo de Binlog via JDBC não é suportado antes da V2.0.18. A partir da V2.0.18, instâncias secundárias suportam o consumo de Binlog, mas não registram o progresso do consumo.

Criar uma publicação e um slot de replicação

Necessário para Hologres V2.0 e anteriores. A partir da V2.1, usuários com apenas permissões de leitura na tabela de destino podem consumir dados do Binlog sem um slot de replicação — pule esta seção se estiver na V2.1 ou posterior.

Publicação

Uma publicação define quais alterações de tabela ficam disponíveis para replicação lógica. No Hologres, uma publicação está vinculada a exatamente uma tabela física, e essa tabela deve ter o Binlog ativado.

Criar uma publicação

CREATE PUBLICATION <name> FOR TABLE <table_name>;

Parâmetro

Descrição

name

Um nome personalizado para a publicação.

table_name

O nome da tabela de destino.

Exemplo:

CREATE PUBLICATION hg_publication_test_1 FOR TABLE test_message_src;

Consultar publicações

SELECT * FROM pg_publication;

Saída de exemplo:

         pubname        | pubowner | puballtables | pubinsert | pubupdate | pubdelete | pubtruncate
------------------------+----------+--------------+-----------+-----------+-----------+-------------
 hg_publication_test_1  |    16728 | f            | t         | t         | t         | t
(1 row)

Campo

Descrição

pubname

O nome da publicação.

pubowner

O proprietário da publicação.

puballtables

Indica se todas as tabelas estão incluídas. Sempre false — não há suporte para vincular múltiplas tabelas físicas.

pubinsert

Indica se eventos INSERT são publicados. Padrão: true.

pubupdate

Indica se eventos UPDATE são publicados. Padrão: true.

pubdelete

Indica se eventos DELETE são publicados. Padrão: true.

pubtruncate

Indica se eventos TRUNCATE são publicados. Padrão: true.

Para mais informações sobre tipos de eventos Binlog, consulte Formato e princípios do Binlog.

Consultar tabelas em uma publicação

SELECT * FROM pg_publication_tables;

Saída de exemplo:

         pubname        | schemaname |    tablename
------------------------+------------+------------------
 hg_publication_test_1  | public     | test_message_src
(1 row)

Excluir uma publicação

DROP PUBLICATION <name>;

Exemplo:

DROP PUBLICATION hg_publication_test_1;

Slot de replicação

Importante

Um slot de replicação inativo continua retendo dados WAL no servidor. Se um slot não for consumido, o WAL se acumula e pode esgotar seu armazenamento. Exclua qualquer slot desnecessário.

Um slot de replicação rastreia o progresso do consumo de uma publicação e permite a retomada da transmissão. Após um failover, um consumidor pode recuperar a partir do último checkpoint confirmado registrado no slot.

Apenas superusuários e usuários com a Replication Role podem criar e usar slots de replicação. Para conceder ou revogar a Replication Role:

-- Grant the Replication Role to a user.
ALTER ROLE <user_name> REPLICATION;

-- Revoke the Replication Role from a user.
ALTER ROLE <user_name> NOREPLICATION;

user_name é um ID de conta Alibaba Cloud ou um usuário do Resource Access Management (RAM). Consulte Visão geral de contas.

Criar um slot de replicação

CALL hg_create_logical_replication_slot('<replication_slot_name>', 'hgoutput', '<publication_name>');

Parâmetro

Descrição

replication_slot_name

Um nome personalizado para o slot de replicação.

hgoutput

O plugin de saída para o formato Binlog. Apenas hgoutput é suportado.

publication_name

O nome da publicação a ser vinculada.

Exemplo:

CALL hg_create_logical_replication_slot('hg_replication_slot_1', 'hgoutput', 'hg_publication_test_1');

Consultar slots de replicação

SELECT * FROM hologres.hg_replication_slot_properties;

Saída de exemplo:

        slot_name       | property_key |      property_value
------------------------+--------------+------------------------
 hg_replication_slot_1  | plugin       | hgoutput
 hg_replication_slot_1  | publication  | hg_publication_test_1
 hg_replication_slot_1  | parallelism  | 1
(3 rows)

Campo

Descrição

slot_name

O nome do slot de replicação.

property_key

Um dos seguintes: plugin (plugin de saída), publication (publicação vinculada) ou parallelism (número de conexões simultâneas necessárias para consumir a tabela inteira, igual à contagem de shards do grupo de tabelas).

property_value

O valor da propriedade correspondente.

Consultar o paralelismo necessário

Como o Hologres é um banco de dados distribuído, os dados de cada tabela são distribuídos por vários shards. É necessária uma conexão por shard para consumir a tabela completa. Para verificar quantas conexões simultâneas o hg_replication_slot_1 requer:

SELECT hg_get_logical_replication_slot_parallelism('hg_replication_slot_1');

Saída de exemplo:

 hg_get_logical_replication_slot_parallelism
---------------------------------------------
                                          20

Consultar o progresso do consumo

A tabela hologres.hg_replication_progress registra o offset do consumidor que você confirma explicitamente.

SELECT * FROM hologres.hg_replication_progress;

Saída de exemplo:

        slot_name       | parallel_index | lsn
------------------------+----------------+-----
 hg_replication_slot_1  |              0 |  66
 hg_replication_slot_1  |              1 | 122
 hg_replication_slot_1  |              2 | 119
(3 rows)

Campo

Descrição

slot_name

O nome do slot de replicação.

parallel_index

O índice da conexão simultânea (uma por shard).

lsn

O Log Sequence Number (LSN) do último registro Binlog consumido que foi confirmado.

Importante
  • A tabela hologres.hg_replication_progress é criada somente após o primeiro consumo de Binlog.

  • A tabela registra apenas offsets confirmados explicitamente pela chamada da função de commit LSN no código. O valor registrado pode não corresponder à posição real do consumidor. Rastreie o LSN no lado do cliente e use-o como ponto de recuperação.

  • Os commits de checkpoint são efetivos apenas ao consumir via slot de replicação. Ao consumir apenas pelo nome da tabela (sem withSlotName), o progresso não é registrado.

Excluir um slot de replicação

Exclua um slot quando ele não for mais necessário para evitar o acúmulo de WAL. Os slots não são fechados automaticamente quando um job de consumo é interrompido — você deve excluí-los manualmente.

CALL hg_drop_logical_replication_slot('<replication_slot_name>');

Exemplo:

CALL hg_drop_logical_replication_slot('hg_replication_slot_1');

Consumir dados do Binlog usando JDBC

O consumo baseado em JDBC oferece controle granular sobre cada shard. Crie um PGReplicationStream por shard e gerencie o paralelismo manualmente.

Etapa 1: Adicionar dependências

Use JDBC 42.2.18 ou posterior. Adicione as seguintes dependências ao seu pom.xml:

<dependency>
    <groupId>org.postgresql</groupId>
    <artifactId>postgresql</artifactId>
    <version>42.3.8</version>
</dependency>
<!-- Used to get the table schema and decode Binlog records -->
<dependency>
    <groupId>com.alibaba.hologres</groupId>
    <artifactId>holo-client</artifactId>
    <version>2.2.10</version>
</dependency>

Etapa 2: Conectar e consumir

O exemplo a seguir conecta-se a um slot de replicação, lê registros Binlog de um shard e imprime cada registro.

import com.alibaba.hologres.client.HoloClient;
import com.alibaba.hologres.client.HoloConfig;
import com.alibaba.hologres.client.impl.binlog.HoloBinlogDecoder;
import com.alibaba.hologres.client.model.Record;
import com.alibaba.hologres.client.model.TableSchema;

import org.postgresql.PGConnection;
import org.postgresql.PGProperty;
import org.postgresql.replication.LogSequenceNumber;
import org.postgresql.replication.PGReplicationStream;

import java.nio.ByteBuffer;
import java.sql.Connection;
import java.sql.DriverManager;
import java.util.Arrays;
import java.util.List;
import java.util.Properties;

public class Test {

    public static void main(String[] args) throws Exception {

        String username = "";
        String password = "";
        String url = "jdbc:postgresql://Endpoint:Port/db_test";

        // Set up the JDBC connection for replication.
        Properties properties = new Properties();
        PGProperty.USER.set(properties, username);
        PGProperty.PASSWORD.set(properties, password);
        PGProperty.ASSUME_MIN_SERVER_VERSION.set(properties, "9.4");
        // Required for Binlog consumption.
        PGProperty.REPLICATION.set(properties, "database");

        try (Connection connection = DriverManager.getConnection(url, properties)) {

            // Each PGReplicationStream consumes one shard. Set shardId for each stream.
            int shardId = 0;
            PGConnection pgConnection = connection.unwrap(PGConnection.class);
            PGReplicationStream pgReplicationStream = pgConnection.getReplicationAPI()
                    .replicationStream()
                    .logical()
                    // V2.1+: two options
                    // Option 1: specify the replication slot name. table_name is ignored.
                    // Option 2: omit withSlotName and set table_name in withSlotOption instead.
                    .withSlotName("slot_name")
                    .withSlotOption("table_name", "public.test_message_src")
                    .withSlotOption("parallel_index", shardId)
                    .withSlotOption("batch_size", "1024")
                    .withSlotOption("start_time", "2021-01-01 00:00:00")
                    .withSlotOption("start_lsn", "0")
                    .start();

            // HoloClient is needed to get the table schema for decoding.
            HoloConfig holoConfig = new HoloConfig();
            holoConfig.setJdbcUrl(url);
            holoConfig.setUsername(username);
            holoConfig.setPassword(password);
            HoloClient client = new HoloClient(holoConfig);

            // Decode raw binary Binlog data into BinlogRecord objects.
            TableSchema schema = client.getTableSchema("test_message_src", true);
            HoloBinlogDecoder decoder = new HoloBinlogDecoder(schema);

            // Track the last consumed LSN for checkpoint recovery.
            Long currentLsn = 0L;

            ByteBuffer byteBuffer = pgReplicationStream.readPending();
            while (true) {
                if (byteBuffer != null) {
                    List<BinlogRecord> records = decoder.decode(shardId, byteBuffer);
                    Long latestLsn = 0L;
                    for (BinlogRecord record : records) {
                        latestLsn = record.getBinlogLsn();
                        // Process the record here.
                        System.out.println("lsn: " + latestLsn + ", record: " + Arrays.toString(record.getValues()));
                    }
                    // Save the latest LSN as the recovery point.
                    currentLsn = latestLsn;
                    pgReplicationStream.forceUpdateStatus();
                }
                byteBuffer = pgReplicationStream.readPending();
            }
        }
    }
}

Parâmetros

Comportamento de withSlotName por versão:

Versão

Comportamento

Anterior à V2.1

Especifique o nome de um slot de replicação existente. Obrigatório.

V2.1 e posteriores

Opcional. Omita withSlotName e defina table_name em withSlotOption para consumir sem um slot de replicação.

Parâmetros de withSlotOption:

Parâmetro

Obrigatório

Descrição

table_name

Sim, quando withSlotName não estiver definido

A tabela de destino a ser consumida. Formato: schema_name.table_name ou table_name. Ignorado se withSlotName estiver definido.

parallel_index

Sim

O índice do shard a ser consumido. Um PGReplicationStream lida com um shard. Para uma tabela com três shards, crie três fluxos com parallel_index 0, 1 e 2.

start_time

Não

Inicia o consumo a partir deste timestamp. Exemplo: 2021-01-01 12:00:00+08.

start_lsn

Não

Inicia o consumo a partir do registro posterior a este LSN. Tem prioridade sobre start_time.

batch_size

Não

Máximo de registros por recuperação. Padrão: 1.024.

Comportamento inicial padrão (quando nem start_lsn nem start_time estão definidos):

Cenário

Comportamento

Primeiro consumo a partir de um slot de replicação

Começa do início, semelhante ao offset earliest do Kafka.

Consumo subsequente a partir de um slot de replicação

Retoma a partir do último checkpoint confirmado.

Consumo apenas pelo nome da tabela (sem withSlotName)

Sempre começa do início.

Campos BinlogRecord

BinlogRecord expõe os seguintes campos de sistema do Binlog:

Método

Retorna

getBinlogLsn()

O LSN do registro Binlog.

getBinlogTimestamp()

O timestamp do sistema quando o evento ocorreu.

getBinlogEventType()

O tipo de evento: INSERT, UPDATE, DELETE ou TRUNCATE.

Para uma referência completa de campos, consulte Assinar dados do Binlog do Hologres.

Confirmar o checkpoint

Após consumir dados do Binlog, confirme o checkpoint chamando a função de commit LSN. Isso serve a dois propósitos: informa ao servidor quais segmentos WAL podem ser arquivados ou descartados com segurança e fornece um ponto de recuperação para que o consumo possa retomar da posição correta após um failover.

O código de exemplo acima não inclui esta etapa — adicione-a com base nos requisitos de recuperação da sua aplicação.

Consumir dados do Binlog usando Holo-Client

O Holo-Client simplifica o consumo de Binlog gerenciando conexões no nível de shard automaticamente. Especifique a tabela de destino e o Holo-Client cuidará do restante.

  • O número de conexões é igual à contagem de shards da tabela.

  • Salve o checkpoint por shard para que o consumo possa retomar da posição correta após uma falha de rede ou outra interrupção.

Nota

Use Holo-Client 2.2.10 ou posterior. As versões 2.2.9 e anteriores apresentam um vazamento de memória.

Etapa 1: Adicionar dependências

<dependency>
    <groupId>com.alibaba.hologres</groupId>
    <artifactId>holo-client</artifactId>
    <version>2.2.10</version>
</dependency>

Etapa 2: Conectar e consumir

O exemplo a seguir assina uma tabela, processa cada registro, salva um checkpoint por shard e tenta novamente em caso de falha.

import com.alibaba.hologres.client.BinlogShardGroupReader;
import com.alibaba.hologres.client.Command;
import com.alibaba.hologres.client.HoloClient;
import com.alibaba.hologres.client.HoloConfig;
import com.alibaba.hologres.client.Subscribe;
import com.alibaba.hologres.client.exception.HoloClientException;
import com.alibaba.hologres.client.impl.binlog.BinlogOffset;
import com.alibaba.hologres.client.model.binlog.BinlogHeartBeatRecord;
import com.alibaba.hologres.client.model.binlog.BinlogRecord;

import java.util.HashMap;
import java.util.Map;

public class HoloBinlogExample {

    public static BinlogShardGroupReader reader;

    public static void main(String[] args) throws Exception {
        String username = "";
        String password = "";
        String url = "jdbc:postgresql://ip:port/database";
        String tableName = "test_message_src";
        String slotName = "hg_replication_slot_1";

        HoloConfig holoConfig = new HoloConfig();
        holoConfig.setJdbcUrl(url);
        holoConfig.setUsername(username);
        holoConfig.setPassword(password);
        holoConfig.setBinlogReadBatchSize(128);
        holoConfig.setBinlogIgnoreDelete(true);
        holoConfig.setBinlogIgnoreBeforeUpdate(true);
        holoConfig.setBinlogHeartBeatIntervalMs(5000L);
        HoloClient client = new HoloClient(holoConfig);

        // Get the shard count to initialize per-shard checkpoint tracking.
        int shardCount = Command.getShardCount(client, client.getTableSchema(tableName));

        // Initialize the checkpoint map with LSN 0 for each shard.
        Map<Integer, Long> shardIdToLsn = new HashMap<>(shardCount);
        for (int i = 0; i < shardCount; i++) {
            shardIdToLsn.put(i, 0L);
        }

        // Before V2.1: tableName and slotName are both required.
        // V2.1 and later: tableName is enough (slotName defaults to "hg_table_name_slot").
        Subscribe subscribe = Subscribe.newStartTimeBuilder(tableName, slotName)
                .setBinlogReadStartTime("2021-01-01 12:00:00")
                .build();

        reader = client.binlogSubscribe(subscribe);

        BinlogRecord record;
        int retryCount = 0;

        while (true) {
            try {
                if (reader.isCanceled()) {
                    // Re-create the reader using the last saved checkpoint.
                    reader = client.binlogSubscribe(subscribe);
                }
                while ((record = reader.getBinlogRecord()) != null) {
                    if (record instanceof BinlogHeartBeatRecord) {
                        // All data up to this heartbeat's timestamp has been consumed on this shard.
                        continue;
                    }

                    // Process the record here.
                    System.out.println(record);

                    // Save the checkpoint so consumption can resume from here on failure.
                    shardIdToLsn.put(record.getShardId(), record.getBinlogLsn());
                    retryCount = 0;
                }
            } catch (HoloClientException e) {
                if (++retryCount > 10) {
                    throw new RuntimeException(e);
                }
                System.out.println(String.format(
                        "Binlog read failed: %s. Retrying (%d/10)...", e.getMessage(), retryCount));

                // Wait before retrying, with increasing backoff.
                Thread.sleep(5000L * retryCount);

                // Resume from the saved per-shard checkpoint.
                Subscribe.OffsetBuilder subscribeBuilder = Subscribe.newOffsetBuilder(tableName, slotName);
                for (int i = 0; i < shardCount; i++) {
                    subscribeBuilder.addShardStartOffset(i,
                            new BinlogOffset().setSequence(shardIdToLsn.get(i)));
                }
                subscribe = subscribeBuilder.build();
                reader.cancel();
            }
        }
    }
}

Parâmetros

Tipos de Subscribe:

Tipo

Quando usar

Subscribe.newStartTimeBuilder(tableName, slotName)

Iniciar o consumo a partir de um timestamp específico.

Subscribe.newOffsetBuilder(tableName, slotName)

Iniciar o consumo a partir de um LSN específico por shard — use isto para recuperação baseada em checkpoint.

Nota sobre versão: Antes da V2.1, tanto tableName quanto slotName são obrigatórios. A partir da V2.1, apenas tableName é necessário (equivalente a usar o nome de slot fixo hg_table_name_slot).

Parâmetros Binlog de HoloConfig:

Parâmetro

Padrão

Descrição

binlogReadBatchSize

1.024

Máximo de registros por recuperação por shard.

binlogHeartBeatIntervalMs

-1 (desativado)

Intervalo em milissegundos entre mensagens BinlogHeartBeatRecord. Quando nenhum dado novo chega, os registros de heartbeat indicam que todos os dados até aquele timestamp foram consumidos naquele shard.

binlogIgnoreDelete

false

Ignorar eventos DELETE.

binlogIgnoreBeforeUpdate

false

Ignorar os registros de imagem anterior (before-image) para eventos UPDATE.

Solução de problemas

A tabela hologres.hg_replication_progress está ausente ou vazia

Se a tabela não existir ou não mostrar dados após você confirmar o progresso do consumo, verifique o seguinte:

Consumo pelo nome da tabela sem slot de replicação. Se você não definiu withSlotName (JDBC) ou não especificou um slot (Holo-Client), o rastreamento de progresso não é suportado. A tabela não será criada nem atualizada. Mude para o consumo baseado em slot de replicação para ativar o rastreamento de progresso.

Primeiro consumo em uma instância secundária somente leitura. Em instâncias secundárias executando Hologres anterior à V2.0.18, a tabela hologres.hg_replication_progress não pode ser criada durante o primeiro consumo de Binlog. Consuma da instância primária uma vez e depois retorne para a instância secundária.

Se nenhuma das situações acima se aplicar, entre em contato com o suporte através do grupo DingTalk do Hologres. Consulte Como obtenho mais suporte online?.

Próximos passos