O conector Hudi integrado não terá mais suporte nas versões futuras do Ververica Runtime (VVR). Use custom connectors para conectar o Realtime Compute for Apache Flink ao Apache Hudi ou migre para o Paimon connector para obter recursos e desempenho otimizados.
O Apache Hudi é um framework de data lake open source que gerencia dados tabulares 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 com viagem no tempo.
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. |
|
Viagem no tempo |
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) — como 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 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 análises OLAP downstream. |
|
Enfileiramento de mensagens |
Funciona como substituto leve para filas de mensagens em cenários de baixo volume e simplifica a arquitetura da aplicação. |
|
Preenchimento retroativo 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
Livre 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: Integra-se ao conector Flink CDC para agilizar o desenvolvimento de dados.
Recursos de classe empresarial: Oferece gerenciamento unificado de metadados via Data Lake Formation (DLF) e alterações leves automáticas de esquema.
Armazenamento econômico: Armazena dados nos formatos Apache Parquet ou Apache Avro no Alibaba Cloud OSS, com isolamento entre 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 das 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 alterar 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 múltiplos 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 Regions and endpoints. Para endpoints do OSS-HDFS, verifique a seção Port na 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 Manage variables. |
|
|
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 sequencialmente em um ou mais buckets. Aumentar este valor não eleva a quantidade de arquivos pequenos. |
|
|
Paralelismo do deployment |
Paralelismo dos operadores de atribuição de bucket. Aumentar este valor incrementa a quantidade de arquivos pequenos. |
|
|
Paralelismo do deployment |
Paralelismo dos operadores de bootstrap de índice. Tem efeito somente 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 roda de forma assíncrona. Defina como |
|
|
|
Paralelismo das tarefas de compactação. |
|
|
|
Estratégia usada para acionar a compactação. Valores válidos: |
|
|
|
Quantidade 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 compactação e deduplicação. Aumente para 1 GB se houver recursos disponíveis. |
|
|
|
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 o tamanho dos 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 do deployment |
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 quantidade 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. Diminui a criação de arquivos pequenos quando uma única tarefa grava em múltiplas 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 habilita 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. Configure compaction.delta_commits=5 e compaction.delta_seconds=3600 para dar tempo suficiente aos consumidores downstream de lerem 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 |
Superior (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 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 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, configurarindex.global.enabled=truenão surte efeito — 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 |
|
Viagem no tempo |
Defina apenas |
Parâmetros de leitura streaming
Por padrão, a leitura de uma tabela Hudi usa leitura de snapshot — o snapshot completo mais recente é retornado de uma 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 intervalo de checkpoint de 5 minutos, o valor padrão 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 afetando potencialmente 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 Settings of DataStream connectors.
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 abaixo 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");
}
}