As futuras versões do Ververica Runtime (VVR) não suportarão mais o conector Hudi integrado. Use conectores personalizados para conectar o Realtime Compute for Apache Flink ao Apache Hudi ou migre para o conector Paimon e obtenha recursos e desempenho otimizados.
O Apache Hudi é um framework de data lake open-source que gerencia dados de tabelas armazenados no Object Storage Service (OSS) ou no Hadoop Distributed File System (HDFS). Ele oferece garantias ACID (Atomicidade, Consistência, Isolamento, Durabilidade), operações de upsert e exclusão em nível de linha, gerenciamento automático de arquivos pequenos e consultas time travel.
Recursos principais
|
Recurso |
Descrição |
|
Semântica ACID |
Oferece isolamento de snapshot por padrão e garante a consistência dos dados durante leituras e gravações simultâneas. |
|
Semântica UPSERT |
Combina INSERT e UPDATE: insere o registro se ele não existir ou o atualiza caso já exista. Isso simplifica o código de desenvolvimento de ETL. |
|
Time travel |
Permite acessar versões históricas dos dados em um ponto específico no tempo, facilitando auditorias eficientes e controle de qualidade. |
Cenários típicos
|
Cenário |
Descrição |
|
Aceleração de ingestão de banco de dados |
Grava dados de Change Data Capture (CDC) (por exemplo, logs binários do MySQL via conector MySQL CDC) diretamente em uma tabela Hudi para ETL em tempo real downstream. Essa abordagem é mais econômica do que o carregamento em massa offline. |
|
ETL incremental |
Extrai fluxos de dados alterados do Hudi de forma incremental para processos de ETL leves e em tempo real. Use Apache Presto ou Apache Spark para OLAP downstream. |
|
Enfileiramento de mensagens |
Substitui filas de mensagens leves em cenários de baixo volume e simplifica a arquitetura da aplicação ao utilizar o Hudi. |
|
Backfilling de dados |
Une dados completos e incrementais de tabelas Hudi em um metastore Hive para gerar tabelas largas com sobrecarga computacional mínima. |
Vantagens sobre o Hudi open-source
Sem necessidade de manutenção: O conector Hudi integrado reduz a complexidade de O&M e fornece garantias de SLA.
Conectividade de dados aprimorada: Desacopla os dados dos mecanismos de computação e permite migração perfeita entre Apache Flink, Apache Spark, Apache Presto e Apache Hive.
Ingestão simplificada de banco para lake: Funciona com o conector Flink CDC para agilizar o desenvolvimento de dados.
Recursos de classe empresarial: Gerenciamento unificado de metadados via Data Lake Formation (DLF) e alterações leves automáticas de esquema.
Armazenamento econômico: Dados armazenados em formato Apache Parquet ou Apache Avro no Alibaba Cloud OSS, com isolamento de armazenamento e computação para dimensionamento flexível de recursos.
Configurações suportadas
|
Item |
Valor |
|
Tipo de tabela |
Tabela source, tabela sink |
|
Modo de execução |
Modo streaming, modo batch |
|
Formato de dados |
N/A |
|
Tipo de API |
DataStream API, SQL API |
|
Atualização/exclusão de dados na sink |
Suportado |
|
Versão mínima do VVR |
vvr-4.0.11-flink-1.13 |
|
Sistemas de arquivos suportados |
OSS, HDFS, OSS-HDFS |
Métricas
|
Tipo de tabela |
Métricas |
|
Tabela source |
|
|
Tabela sink |
|
Para definições de métricas, consulte Métricas.
Limitações
Versão mínima do mecanismo: vvr-4.0.11-flink-1.13 ou posterior.
Sistemas de arquivos suportados: apenas OSS, HDFS ou OSS-HDFS.
Jobs rascunho não podem ser executados em clusters de sessão.
O conector Hudi não suporta modificações de campos. Para modificar campos, use instruções Spark SQL no console do Data Lake Formation (DLF).
Sintaxe
CREATE TEMPORARY TABLE hudi_tbl (
uuid BIGINT,
data STRING,
ts TIMESTAMP(3),
PRIMARY KEY(uuid) NOT ENFORCED
) WITH (
'connector' = 'hudi',
'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
...
);
Parâmetros na cláusula WITH
Parâmetros básicos
Parâmetros comuns
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Sim |
— |
Defina como |
|
|
Sim |
— |
Caminho de armazenamento da tabela. Formatos suportados: OSS ( |
|
|
Não |
|
Campo de chave primária. Separe vários campos com vírgulas. Alternativamente, use a sintaxe |
|
|
Não |
|
Campo de versão usado para determinar a ordem das atualizações. Se não definido, as atualizações seguem a sequência de mensagens definida pelo mecanismo. |
|
|
Não |
— |
Obrigatório ao armazenar dados no OSS ou OSS-HDFS. Para endpoints do OSS, consulte Regiões e endpoints. Para endpoints OSS-HDFS, consulte a seção Port da página Overview do bucket OSS. |
|
|
Não |
— |
AccessKey ID. Obrigatório para OSS e OSS-HDFS. Armazene credenciais como variáveis em vez de codificá-las diretamente. Consulte Gerenciar variáveis. |
|
|
Não |
— |
AccessKey secret. Obrigatório para OSS e OSS-HDFS. |
Para proteger seu par de AccessKey, armazene o AccessKey ID e o AccessKey secret como variáveis. Consulte Gerenciar variáveis.
Parâmetros da tabela source
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Não |
|
Defina como |
|
|
Não |
(vazio) |
Offset inicial para leitura em streaming. Formato: |
Parâmetros da tabela sink
|
Parâmetro |
Obrigatório |
Padrão |
Descrição |
|
|
Não |
|
Modo de gravação. Valores válidos: |
|
|
Não |
|
Defina como |
|
|
Não |
|
Modo de sincronização. |
|
|
Não |
|
Nome do banco de dados Hive de destino. |
|
|
Não |
Nome da tabela atual |
Nome da tabela Hive de destino. Não deve conter hifens ( |
|
|
Não |
— |
Região onde o DLF está ativado. Tem efeito apenas quando |
|
|
Não |
— |
Endpoint do DLF. Tem efeito apenas quando |
Parâmetros avançados
Parâmetros de paralelismo
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Paralelismo das tarefas de gravação. Cada tarefa grava em um ou mais buckets sequencialmente. Aumentar este valor não incrementa a quantidade de arquivos pequenos. |
|
|
Paralelismo da implantação |
Paralelismo dos operadores de atribuição de bucket. Aumentar este valor eleva a contagem de arquivos pequenos. |
|
|
Paralelismo da implantação |
Paralelismo dos operadores de bootstrap de índice. Tem efeito apenas quando |
|
|
|
Paralelismo dos operadores de leitura em streaming e batch. |
|
|
|
Paralelismo dos operadores de compactação online. A compactação online consome mais recursos que a offline; prefira a compactação offline para cargas de trabalho em produção. |
Parâmetros de compactação online
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Define se os planos de compactação são gerados conforme agendamento. Mantenha como |
|
|
|
Define se a compactação é executada de forma assíncrona. Defina como |
|
|
|
Paralelismo das tarefas de compactação. |
|
|
|
Estratégia usada para acionar a compactação. Valores válidos: |
|
|
|
Número de commits necessários para acionar a compactação. Usado com |
|
|
|
Intervalo em segundos entre os acionamentos de compactação. Usado com |
|
|
|
Memória máxima para o mapa hash usado durante a compactação e deduplicação. Aumente para 1 GB se os recursos permitirem. |
|
|
|
Throughput máximo de I/O por plano de compactação. |
Parâmetros de tamanho de arquivo
Estes parâmetros controlam como o Hudi gerencia tamanhos de arquivos para evitar o acúmulo de arquivos pequenos.
|
Parâmetro |
Padrão |
Descrição |
|
|
120 MB (120 × 1024 × 1024 bytes) |
Tamanho máximo de um arquivo Parquet. Dados que excedem esse limite são gravados em um novo grupo de arquivos. |
|
|
100 MB (104.857.600 bytes) |
Arquivos menores que este limite são tratados como arquivos pequenos. Durante gravações, o Hudi anexa dados aos arquivos pequenos existentes em vez de criar novos. |
|
|
1 KB (1.024 bytes) |
Tamanho estimado do registro. Se não definido, o Hudi calcula isso dinamicamente a partir dos metadados commitados. |
Parâmetros de configuração do Hadoop
|
Parâmetro |
Padrão |
Descrição |
|
|
— |
Itens de configuração do Hadoop, especificados com o prefixo |
Parâmetros de gravação de dados
Gravação em lote
Use a gravação em lote para importar dados existentes de outras fontes para uma tabela Hudi.
bulk_insertignora serialização Avro, compactação e deduplicação. Garanta a unicidade dos dados de origem antes de usar este modo.bulk_inserté válido apenas no modo de execução em lote.
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Tipo de gravação. Defina como |
|
|
Paralelismo da implantação |
Paralelismo para tarefas |
|
|
|
Define se os dados de entrada devem ser reorganizados (shuffle) pelo campo de partição antes da gravação. Disponível no Hudi 0.11.0 e posteriores. Reduz a contagem de arquivos pequenos, mas pode causar skew de dados. |
|
|
|
Define se os dados de entrada devem ser ordenados pelo campo de partição antes da gravação. Disponível no Hudi 0.11.0 e posteriores. Reduz a contagem de arquivos pequenos quando uma única tarefa grava em várias partições. |
|
|
|
Memória gerenciada disponível para o operador de ordenação, em MB. |
Modo changelog
No modo changelog, o Hudi retém todos os eventos de alteração (INSERT, UPDATE_BEFORE, UPDATE_AFTER e DELETE) e permite data warehousing quase em tempo real ponta a ponta com computação stateful do Flink. Tabelas Merge On Read (MOR) suportam este modo.
No modo sem changelog, alterações intermediárias dentro de um lote são mescladas. A leitura de snapshot retorna apenas o resultado final mesclado; estados intermediários não ficam visíveis, independentemente do caminho de gravação.
Após ativar o modo changelog, a tarefa de compactação assíncrona ainda mescla alterações intermediárias. Defina compaction.delta_commits=5 e compaction.delta_seconds=3600 para dar aos consumidores downstream tempo suficiente para ler os registros antes da compactação.
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Defina como |
Modo append
Suportado no Hudi 0.10.0 e posteriores.
Tabelas MOR: Aplica-se a política de arquivos pequenos. Os dados são gravados em arquivos de log Apache Avro no modo append.
Tabelas Copy On Write (COW): A política de arquivos pequenos não se aplica. Um novo arquivo Apache Parquet é criado para cada gravação.
Parâmetros de clustering
O Hudi suporta clustering para resolver o acúmulo de arquivos pequenos no modo INSERT.
Clustering inline (apenas tabelas COW)
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Defina como |
Clustering assíncrono (Hudi 0.12.0 e posteriores)
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Defina como |
|
|
|
Número de commits necessários para gerar um plano de clustering. Tem efeito apenas quando |
|
|
|
Defina como |
|
|
|
Paralelismo das tarefas de clustering. |
|
|
1 GiB (1.073.741.824 bytes) |
Tamanho máximo alvo do arquivo para saída do clustering. |
|
|
|
Arquivos menores que este limite (em bytes) são elegíveis para clustering. |
|
|
— |
Colunas usadas para ordenar dados durante o clustering. |
Estratégias de plano de clustering
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Modo de filtro de partição. Valores válidos: |
|
|
|
Número de dias recentes para selecionar partições para clustering. Tem efeito apenas quando |
|
|
— |
Partição inicial para filtragem por intervalo. Tem efeito apenas quando |
|
|
— |
Partição final para filtragem por intervalo. Tem efeito apenas quando |
|
|
— |
Expressão regular para seleção de partições. |
|
|
— |
Lista separada por vírgulas das partições selecionadas. |
Escolha um tipo de índice
O Hudi suporta dois tipos de índice. Use esta tabela para escolher aquele adequado à sua carga de trabalho.
|
Dimensão |
FLINK_STATE |
BUCKET |
|
Sobrecarga de armazenamento/computação |
Sim (state backend) |
Nenhuma |
|
Desempenho |
Depende do state backend |
Melhor (sem sobrecarga de estado) |
|
Flexibilidade de grupos de arquivos |
Atribui registros dinamicamente com base no tamanho do arquivo |
Número fixo de buckets (não pode aumentar após a configuração inicial) |
|
Alterações entre partições |
Suportado |
Não suportado (exceção: entrada streaming de Change Data Capture (CDC)) |
|
Quando usar |
Tabelas com menos de 500 milhões de registros ou cargas de trabalho que exigem atualizações entre partições |
Tabelas com mais de 500 milhões de registros onde a sobrecarga de estado é um gargalo |
Quandoindex.typeestá definido comoBUCKET, definirindex.global.enabled=truenão tem efeito, pois o índice bucket não suporta deduplicação entre partições.
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Tipo de índice. Valores válidos: |
|
|
Chave primária |
Campo de chave hash para índice bucket. Pode ser um subconjunto da chave primária. |
|
|
|
Número de buckets por partição. Não pode ser alterado após a criação da tabela. |
Os parâmetros de índice bucket são suportados no Hudi 0.11.0 e posteriores.
Parâmetros de leitura de dados
O Hudi suporta três padrões de leitura usando o mesmo conjunto de parâmetros.
|
Padrão |
Configuração |
|
Leitura streaming |
Defina |
|
Leitura batch incremental |
Defina tanto |
|
Time travel |
Defina apenas |
Parâmetros de leitura streaming
Por padrão, a leitura de uma tabela Hudi usa leitura de snapshot, retornando o snapshot completo mais recente de uma só vez. Defina read.streaming.enabled=true para alternar para leitura streaming.
|
Parâmetro |
Padrão |
Descrição |
|
|
|
Defina como |
|
|
(vazio) |
Offset inicial. Formato: |
|
|
|
Número máximo de commits históricos retidos pelo limpador. Commits além desse limite são excluídos. Por exemplo, com um intervalo de checkpoint de 5 minutos, o valor padrão de 30 retém changelogs por pelo menos 150 minutos. |
A leitura streaming de changelogs requer Hudi 0.10.0 ou posterior. Tarefas de compactação podem mesclar changelogs, removendo registros intermediários e potencialmente afetando cálculos downstream.
Parâmetros de leitura incremental
|
Parâmetro |
Padrão |
Descrição |
|
|
Commit mais recente |
Início do intervalo de leitura, no formato |
|
|
Commit mais recente |
Fim do intervalo de leitura, no formato |
Exemplos
Tabela source
CREATE TEMPORARY TABLE blackhole (
id INT NOT NULL PRIMARY KEY NOT ENFORCED,
data STRING,
ts TIMESTAMP(3)
) WITH (
'connector' = 'blackhole'
);
CREATE TEMPORARY TABLE hudi_tbl (
id INT NOT NULL PRIMARY KEY NOT ENFORCED,
data STRING,
ts TIMESTAMP(3)
) WITH (
'connector' = 'hudi',
'oss.endpoint' = '<yourOSSEndpoint>',
'accessKeyId' = '${secret_values.ak_id}',
'accessKeySecret' = '${secret_values.ak_secret}',
'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
'table.type' = 'MERGE_ON_READ',
'read.streaming.enabled' = 'true'
);
-- Read from the latest commit in streaming mode and write to Blackhole.
INSERT INTO blackhole SELECT * FROM hudi_tbl;
Tabela sink
CREATE TEMPORARY TABLE datagen (
id INT NOT NULL PRIMARY KEY NOT ENFORCED,
data STRING,
ts TIMESTAMP(3)
) WITH (
'connector' = 'datagen',
'rows-per-second' = '100'
);
CREATE TEMPORARY TABLE hudi_tbl (
id INT NOT NULL PRIMARY KEY NOT ENFORCED,
data STRING,
ts TIMESTAMP(3)
) WITH (
'connector' = 'hudi',
'oss.endpoint' = '<yourOSSEndpoint>',
'accessKeyId' = '${secret_values.ak_id}',
'accessKeySecret' = '${secret_values.ak_secret}',
'path' = 'oss://<yourOSSBucket>/<Custom storage directory>',
'table.type' = 'MERGE_ON_READ'
);
INSERT INTO hudi_tbl SELECT * FROM datagen;
DataStream API
Para usar a DataStream API, configure um conector DataStream para o Realtime Compute for Apache Flink. Consulte Configurações de conectores DataStream.
Dependências Maven
Alinhe as versões das dependências com sua versão do VVR.
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<flink.version>1.15.4</flink.version>
<hudi.version>0.13.1</hudi.version>
</properties>
<dependencies>
<!-- Flink -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-common</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner_2.12</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<!-- Hudi -->
<dependency>
<groupId>org.apache.hudi</groupId>
<artifactId>hudi-flink1.15-bundle</artifactId>
<version>${hudi.version}</version>
<scope>provided</scope>
</dependency>
<!-- OSS -->
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>3.3.2</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-aliyun</artifactId>
<version>3.3.2</version>
<scope>provided</scope>
</dependency>
<!-- DLF -->
<dependency>
<groupId>com.aliyun.datalake</groupId>
<artifactId>metastore-client-hive2</artifactId>
<version>0.2.14</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-mapreduce-client-core</artifactId>
<version>2.5.1</version>
<scope>provided</scope>
</dependency>
</dependencies>
As dependências do DLF entram em conflito com versões open-source do Apache Hive (hive-common, hive-exec). Para testes locais com DLF, baixe os pacotes JAR personalizados hive-common e hive-exec e importe-os manualmente no IntelliJ IDEA.
Gravar dados no Hudi
O exemplo a seguir grava dados em uma tabela Hudi MOR no OSS e, opcionalmente, sincroniza metadados com o DLF.
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.data.GenericRowData;
import org.apache.flink.table.data.RowData;
import org.apache.flink.table.data.StringData;
import org.apache.hudi.common.model.HoodieTableType;
import org.apache.hudi.configuration.FlinkOptions;
import org.apache.hudi.util.HoodiePipeline;
import java.util.HashMap;
import java.util.Map;
public class FlinkHudiQuickStart {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
String dbName = "test_db";
String tableName = "test_tbl";
String basePath = "oss://xxx";
Map<String, String> options = new HashMap<>();
// Hudi configuration
options.put(FlinkOptions.PATH.key(), basePath);
options.put(FlinkOptions.TABLE_TYPE.key(), HoodieTableType.MERGE_ON_READ.name());
options.put(FlinkOptions.PRECOMBINE_FIELD.key(), "ts");
options.put(FlinkOptions.DATABASE_NAME.key(), dbName);
options.put(FlinkOptions.TABLE_NAME.key(), tableName);
// OSS configuration
// Use the public endpoint for local debugging (e.g., oss-cn-hangzhou.aliyuncs.com)
// Use the internal endpoint for cluster submission (e.g., oss-cn-hangzhou-internal.aliyuncs.com)
options.put("hadoop.fs.oss.accessKeyId", "xxx");
options.put("hadoop.fs.oss.accessKeySecret", "xxx");
options.put("hadoop.fs.oss.endpoint", "xxx");
options.put("hadoop.fs.AbstractFileSystem.oss.impl", "org.apache.hadoop.fs.aliyun.oss.OSS");
options.put("hadoop.fs.oss.impl", "org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem");
// DLF configuration (optional — remove if not syncing to DLF)
// Use the public endpoint for local debugging (e.g., dlf.cn-hangzhou.aliyuncs.com)
// Use the VPC endpoint for cluster submission (e.g., dlf-vpc.cn-hangzhou.aliyuncs.com)
options.put(FlinkOptions.HIVE_SYNC_ENABLED.key(), "true");
options.put(FlinkOptions.HIVE_SYNC_MODE.key(), "hms");
options.put(FlinkOptions.HIVE_SYNC_DB.key(), dbName);
options.put(FlinkOptions.HIVE_SYNC_TABLE.key(), tableName);
options.put("hadoop.dlf.catalog.id", "xxx");
options.put("hadoop.dlf.catalog.accessKeyId", "xxx");
options.put("hadoop.dlf.catalog.accessKeySecret", "xxx");
options.put("hadoop.dlf.catalog.region", "xxx");
options.put("hadoop.dlf.catalog.endpoint", "xxx");
options.put("hadoop.hive.imetastoreclient.factory.class",
"com.aliyun.datalake.metastore.hive2.DlfMetaStoreClientFactory");
DataStream<RowData> dataStream = env.fromElements(
GenericRowData.of(StringData.fromString("id1"), StringData.fromString("name1"), 22,
StringData.fromString("1001"), StringData.fromString("p1")),
GenericRowData.of(StringData.fromString("id2"), StringData.fromString("name2"), 32,
StringData.fromString("1002"), StringData.fromString("p2"))
);
HoodiePipeline.Builder builder = HoodiePipeline.builder(tableName)
.column("uuid string")
.column("name string")
.column("age int")
.column("ts string")
.column("`partition` string")
.pk("uuid")
.partition("partition")
.options(options);
// Second parameter: whether the input stream is bounded (true = batch, false = streaming)
builder.sink(dataStream, false);
env.execute("Flink_Hudi_Quick_Start");
}
}