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_binlogcriada (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.ImportanteNunca use
DROP EXTENSION <extension_name> CASCADE. A opçãoCASCADEremove 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 DATABASEe a Replication Role na instância
NotaA 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:
Identifique jobs ativos de consumo de Binlog via JDBC e interrompa os desnecessários.
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.
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 countExemplo: 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 alreadyPara 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 |
|
|
Um nome personalizado para a publicação. |
|
|
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 |
|
|
O nome da publicação. |
|
|
O proprietário da publicação. |
|
|
Indica se todas as tabelas estão incluídas. Sempre |
|
|
Indica se eventos INSERT são publicados. Padrão: |
|
|
Indica se eventos UPDATE são publicados. Padrão: |
|
|
Indica se eventos DELETE são publicados. Padrão: |
|
|
Indica se eventos TRUNCATE são publicados. Padrão: |
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
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 |
|
|
Um nome personalizado para o slot de replicação. |
|
|
O plugin de saída para o formato Binlog. Apenas |
|
|
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 |
|
|
O nome do slot de replicação. |
|
|
Um dos seguintes: |
|
|
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 |
|
|
O nome do slot de replicação. |
|
|
O índice da conexão simultânea (uma por shard). |
|
|
O Log Sequence Number (LSN) do último registro Binlog consumido que foi confirmado. |
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 |
Parâmetros de withSlotOption:
|
Parâmetro |
Obrigatório |
Descrição |
|
|
Sim, quando |
A tabela de destino a ser consumida. Formato: |
|
|
Sim |
O índice do shard a ser consumido. Um |
|
|
Não |
Inicia o consumo a partir deste timestamp. Exemplo: |
|
|
Não |
Inicia o consumo a partir do registro posterior a este LSN. Tem prioridade sobre |
|
|
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 |
Sempre começa do início. |
Campos BinlogRecord
BinlogRecord expõe os seguintes campos de sistema do Binlog:
|
Método |
Retorna |
|
|
O LSN do registro Binlog. |
|
|
O timestamp do sistema quando o evento ocorreu. |
|
|
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.
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 |
|
|
Iniciar o consumo a partir de um timestamp específico. |
|
|
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 |
|
|
1.024 |
Máximo de registros por recuperação por shard. |
|
|
-1 (desativado) |
Intervalo em milissegundos entre mensagens |
|
|
|
Ignorar eventos DELETE. |
|
|
|
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
Assinar dados do Binlog do Hologres — Formato Binlog, tipos de eventos e campos de sistema
Melhores práticas para configurar grupos de tabelas — Otimize a contagem de shards para controlar o uso de Walsender
Gerenciamento de instâncias — Contagens de nós frontend por especificação de instância