O plugin wal2json decodifica alterações do write-ahead log (WAL) para o formato JSON, facilitando o consumo de eventos de alteração de banco de dados em pipelines de dados, sistemas de captura de dados alterados (CDC), ferramentas de auditoria e serviços orientados a eventos.
Pré-requisitos
Antes de começar, verifique se você tem:
Um cluster PolarDB for PostgreSQL executando uma versão secundária do mecanismo compatível
O parâmetro
wal_leveldefinido comological. Defina essa configuração no console em Set cluster parameters. Essa alteração exige a reinicialização do cluster.
Compatibilidade
As seguintes versões secundárias do mecanismo são compatíveis com o wal2json:
|
Versão do PostgreSQL |
Versão secundária mínima do mecanismo |
|
PostgreSQL 18 |
2.0.18.1.1.0 |
|
PostgreSQL 17 |
2.0.17.7.5.0 |
|
PostgreSQL 16 |
2.0.16.6.2.0 |
|
PostgreSQL 15 |
2.0.15.12.4.0 |
|
PostgreSQL 14 |
2.0.14.5.1.0 |
|
PostgreSQL 11 |
2.0.11.9.29.0 |
Para verificar sua versão, execute SHOW polardb_version; ou visualize-a no console. Caso sua versão não seja compatível, atualize a versão secundária do mecanismo.
Funcionamento
O wal2json é um plugin de saída de decodificação lógica para PostgreSQL. Em vez de usar o formato padrão pgoutput, ele converte entradas WAL em JSON, simplificando a integração das alterações WAL com sistemas que consomem dados nesse formato.
Para cada transação, o wal2json gera um objeto JSON contendo todas as tuplas (linhas) alteradas. É possível incluir metadados adicionais na saída — como IDs de transação, carimbos de data/hora, esquemas e tipos de dados — utilizando os parâmetros descritos abaixo.
O plugin oferece duas formas de consumir alterações:
API SQL — Use
pg_logical_slot_get_changes()para obter alterações sob demanda. Essa abordagem é ideal para processamento em lote e consultas pontuais.Protocolo de streaming — Use slots de replicação lógica com um consumidor de streaming (por exemplo,
pg_recvlogical). Essa opção é recomendada para pipelines de CDC em tempo real.
O plugin lida com as seguintes operações DML: INSERT, UPDATE, DELETE e TRUNCATE (apenas na versão 2 do formato).
Observações de uso
Registro completo de linhas por padrão
O PolarDB for PostgreSQL define REPLICA_IDENTITY_FULL em todas as tabelas por padrão. Isso significa que cada operação UPDATE e DELETE grava a linha completa no WAL, e não apenas as colunas alteradas. Para registrar somente as colunas modificadas, desative o parâmetro polar_create_table_with_full_replica_identity. Não é possível alterar esse parâmetro pelo console. Entre em contato conosco para obter assistência.
Tabelas sem chave primária
Em tabelas sem chave primária, o wal2json utiliza REPLICA_IDENTITY para determinar quais informações da linha serão incluídas na saída de UPDATE e DELETE. Com a configuração padrão REPLICA_IDENTITY_FULL do PolarDB for PostgreSQL, os dados completos da linha são registrados para todas as tabelas, inclusive aquelas sem chave primária.
Recuperar alterações usando SQL
Não instale o wal2json via CREATE EXTENSION. Carregue-o por meio de um slot de replicação lógica.
O exemplo a seguir utiliza uma tabela orders para demonstrar como o wal2json captura operações INSERT e UPDATE:
-
Crie a tabela e um slot de replicação lógica:
-- Create a sample table CREATE TABLE orders ( order_id SERIAL PRIMARY KEY, product_name VARCHAR(100), quantity INT, status TEXT ); -- Create a logical replication slot using the wal2json plugin SELECT 'init' FROM pg_create_logical_replication_slot('orders_slot', 'wal2json'); -
Faça alterações na tabela:
BEGIN; INSERT INTO orders (product_name, quantity, status) VALUES ('Laptop', 2, 'pending'); INSERT INTO orders (product_name, quantity, status) VALUES ('Mouse', 5, 'pending'); UPDATE orders SET status = 'confirmed' WHERE product_name = 'Laptop'; COMMIT; -
Recupere as alterações como JSON:
SELECT data FROM pg_logical_slot_get_changes('orders_slot', NULL, NULL, 'pretty-print', '1');A saída terá a seguinte aparência:
{ "change": [ { "kind": "insert", "schema": "public", "table": "orders", "columnnames": ["order_id", "product_name", "quantity", "status"], "columntypes": ["integer", "character varying(100)", "integer", "text"], "columnvalues": [1, "Laptop", 2, "pending"] }, { "kind": "insert", "schema": "public", "table": "orders", "columnnames": ["order_id", "product_name", "quantity", "status"], "columntypes": ["integer", "character varying(100)", "integer", "text"], "columnvalues": [2, "Mouse", 5, "pending"] }, { "kind": "update", "schema": "public", "table": "orders", "columnnames": ["order_id", "product_name", "quantity", "status"], "columntypes": ["integer", "character varying(100)", "integer", "text"], "columnvalues": [1, "Laptop", 2, "confirmed"], "oldkeys": { "keynames": ["order_id"], "keytypes": ["integer"], "keyvalues": [1] } } ] } -
Ao concluir, exclua o slot de replicação:
-- Returns 'stop' on success SELECT 'stop' FROM pg_drop_replication_slot('orders_slot');
Parâmetros
Passe todos os parâmetros como pares chave-valor para pg_logical_slot_get_changes(). Por exemplo:
SELECT data FROM pg_logical_slot_get_changes(
'my_slot', NULL, NULL,
'pretty-print', '1',
'include-xids', '1'
);
Conteúdo da saída
|
Parâmetro |
Padrão |
Descrição |
|
|
false |
Adiciona o ID da transação ( |
|
|
false |
Inclui um carimbo de data/hora em cada conjunto de alterações |
|
|
false |
Acrescenta o próximo LSN ( |
|
|
true |
Adiciona o nome do esquema a cada alteração |
|
|
true |
Inclui os tipos de dados em cada alteração |
|
|
true |
Adiciona modificadores de tipo — por exemplo, |
|
|
false |
Inclui OIDs de tipo em cada alteração |
|
|
false |
Adiciona informações da restrição |
|
|
false |
Formata o JSON com espaços em branco e recuo |
|
|
false |
Emite a saída após cada alteração individual, em vez de aguardar o conjunto completo de alterações |
Filtragem
|
Parâmetro |
Padrão |
Descrição |
|
|
(nenhum) |
Exclui tabelas específicas. Separe múltiplas entradas com vírgulas e inclua o nome do esquema (por exemplo, |
|
|
todas as tabelas |
Decodifica apenas as tabelas especificadas. A sintaxe é a mesma de |
|
|
(nenhum) |
Exclui linhas com os prefixos de mensagem especificados. Separe múltiplos prefixos com vírgulas. Geralmente usado com |
|
|
todos |
Inclui apenas linhas com os prefixos de mensagem especificados. Aplique |
Formato e operações
|
Parâmetro |
Padrão |
Descrição |
|
|
1 |
Versão do formato de saída. A Versão 1 produz um objeto JSON por transação — indicada para processamento em lote. A Versão 2 gera um objeto JSON por linha, com marcadores opcionais de transação ( |
|
|
todas |
Operações a serem incluídas na saída: |
Glossário de parâmetros do wal2json
|
Termo |
Descrição |
|
|
Uma única entrada WAL referente a uma operação DML (INSERT, UPDATE, DELETE ou TRUNCATE) |
|
|
Um conjunto de entradas |
Exemplo: uso de include-xids
Este exemplo demonstra como passar parâmetros para recuperar IDs de transação junto com os dados de alteração.
-
Crie uma tabela, um slot de replicação e insira uma linha:
DROP TABLE IF EXISTS tbl; CREATE TABLE tbl (id int); SELECT 'init' FROM pg_create_logical_replication_slot('regression_slot', 'wal2json'); INSERT INTO tbl VALUES (1); -
Consulte o slot com
include-xidsativado e verifique se aparece exatamente um ID de transação distinto:SELECT count(*) = 1, count(distinct ((data::json)->'xid')::text) = 1 FROM pg_logical_slot_get_changes( 'regression_slot', NULL, NULL, 'format-version', '1', 'include-xids', '1' );
Próximos passos
Set cluster parameters — ative
wal_level = logicalwal2json no GitHub — documentação oficial e princípios de design