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 |
|
|
Tipos de API |
DataStream e SQL |
|
Suporte a atualizações ou exclusões em tabelas sink |
Sim |
Recursos
|
Recurso |
Descrição |
|
Permite ler dados do Hologres com ou sem binlog, tanto no modo CDC quanto no modo não-CDC. |
|
|
Oferece suporte a consumo completo, incremental e unificado (completo e incremental). |
|
|
Possibilita ignorar novos dados, substituir linhas inteiras ou atualizar apenas campos específicos. |
|
|
Atualiza somente as colunas modificadas em vez de reescrever a linha inteira. |
|
|
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. |
|
|
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:
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_STREAMouCOPY_BULK_LOAD, os valores das colunas da tabela que usamCURRENT_TIMESTAMPouNOW()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 binloghg_binlog_timestamp_usfornece 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.
ImportanteCaso 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
holohubfoi descontinuado no Hologres 2,0 e posteriores. O sistema troca automaticamente para o modojdbc, 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âmetroconnection.fixed.enabledé definido comotrue). -
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.
ImportanteNo 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.
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 parajdbc_fixede 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 parajdbc_fixede definirá 'deduplication.enabled'='false' para evitar deduplicação.ImportanteO service
rpcfoi descontinuado no Hologres 2,0 e posteriores. Se você definir este parâmetro comorpc, o Flink alterará automaticamente o valor parajdbc_fixed. Se você definir o parâmetro com um valor diferente, o Flink usará o valor especificado.O modo
rpcfoi removido no VVR 11,1 e posteriores. Recomendamos o uso do modojdbcpara conexões.Para operações de gravação em cenários de alta concorrência, recomendamos o uso do modo
jdbc_copyouCOPY_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
rpcfoi descontinuado para instâncias do Hologres versão 2,0 ou posterior. Se você definir este parâmetro comorpc, o Flink alterará automaticamente o valor parajdbc_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
ImportanteO modo
rpcfoi removido no VVR 11,1 e posteriores. Por padrão, o modojdbcé usado para conexões. Você também pode ativar o modo de conexão leve definindo o parâmetroconnection.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
UPDATEgera dois registros de binlog consecutivos: o registroupdate_beforepara os dados antigos, seguido pelo registroupdate_afterpara os novos dados.Evite executar
TRUNCATEou 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
DECIMALseja consistente entre o Flink e o Hologres para evitar erros. Para mais informações, consulte FAQ.Ao usar o modo
initialpara 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.
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+
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.
);
Defina
source.binlog.startup-modecomoINITIALpara realizar uma leitura inicial completa da tabela antes de alternar para o consumo incremental do Binlog.Se o parâmetro
startTimeestiver definido, seja diretamente ou selecionando um horário de início na interface de inicialização, ele terá precedência e definirábinlogStartUpModecomo modotimestamp, substituindo outras configurações de modo, pois o parâmetrostartTimetem 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.
);
Defina
binlogStartUpModecomoinitialpara realizar uma leitura inicial completa da tabela antes de alternar para o consumo incremental do Binlog.Se o parâmetro
startTimeestiver definido, seja diretamente ou selecionando um horário de início na interface de inicialização, ele terá precedência e definirábinlogStartUpModecomo modotimestamp, substituindo outras configurações de modo, pois o parâmetrostartTimetem 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.
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.
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
partitionRoutercomotrue.Para permitir que o conector crie automaticamente uma partição filha caso ela não exista, defina
createparttablecomotrue.
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;
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.ImportantePara 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_precreatedeve 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
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+
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();
}
}
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:
|
|
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
Para saber como criar e usar catálogos do Hologres, consulte Manage Hologres catalogs.
O Hologres se integra ao Flink para fornecer uma solução unificada na construção de um data warehouse em tempo real. Para detalhes, consulte Build a Hologres real-time data warehouse.
O Hologres suporta atualizações e correções eficientes de dados, tornando-o adequado para a construção de tabelas wide em cenários de gravação multi-stream. Para um exemplo, consulte User behavior analysis with Flink, MongoDB, and Hologres.