Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Hudi (retiring)

Última atualização: Jun 27, 2026
Importante

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

numRecordsIn, numRecordsInPerSecond

Tabela sink

numRecordsOut, numRecordsOutPerSecond, currentSendTime

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

connector

Sim

Defina como hudi.

path

Sim

Caminho de armazenamento da tabela. Formatos suportados: OSS (oss://<bucket>/<user-defined-dir>), HDFS (hdfs://<user-defined-dir>), OSS-HDFS (oss://<bucket>.<oss-hdfs-endpoint>/<user-defined-dir>). Caminhos OSS-HDFS exigem VVR 8.0.3 ou posterior. Encontre o endpoint OSS-HDFS na seção Port da página Overview do bucket OSS.

hoodie.datasource.write.recordkey.field

Não

uuid

Campo de chave primária. Separe vários campos com vírgulas. Alternativamente, use a sintaxe PRIMARY KEY na DDL.

precombine.field

Não

ts

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.

oss.endpoint

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.

accessKeyId

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.

accessKeySecret

Não

AccessKey secret. Obrigatório para OSS e OSS-HDFS.

Importante

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

read.streaming.enabled

Não

false

Defina como true para ativar a leitura em streaming. Por padrão, o sistema usa a leitura de snapshot, que retorna o snapshot completo mais recente.

read.start-commit

Não

(vazio)

Offset inicial para leitura em streaming. Formato: yyyyMMddHHmmss para um horário específico ou earliest para ler desde o início. Deixe em branco para ler a partir do commit mais recente.

Parâmetros da tabela sink

Parâmetro

Obrigatório

Padrão

Descrição

write.operation

Não

UPSERT

Modo de gravação. Valores válidos: insert (anexar), upsert (inserir ou atualizar), bulk_insert (anexar em lote).

hive_sync.enable

Não

false

Defina como true para sincronizar metadados com o Apache Hive.

hive_sync.mode

Não

hms

Modo de sincronização. hms sincroniza com um metastore Hive ou DLF. jdbc sincroniza via driver Java Database Connectivity (JDBC).

hive_sync.db

Não

default

Nome do banco de dados Hive de destino.

hive_sync.table

Não

Nome da tabela atual

Nome da tabela Hive de destino. Não deve conter hifens (-).

dlf.catalog.region

Não

Região onde o DLF está ativado. Tem efeito apenas quando hive_sync.mode é hms. Consulte Regiões e endpoints suportados. Deve corresponder à região especificada por dlf.catalog.endpoint.

dlf.catalog.endpoint

Não

Endpoint do DLF. Tem efeito apenas quando hive_sync.mode é hms. Use o endpoint VPC para menor latência (por exemplo, dlf-vpc.cn-hangzhou.aliyuncs.com para a região China (Hangzhou)). Consulte Regiões e endpoints suportados. Para acesso entre VPCs, consulte Como o Realtime Compute for Apache Flink acessa um serviço entre VPCs?

Parâmetros avançados

Parâmetros de paralelismo

Parâmetro

Padrão

Descrição

write.tasks

4

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.

write.bucket_assign.tasks

Paralelismo da implantação

Paralelismo dos operadores de atribuição de bucket. Aumentar este valor eleva a contagem de arquivos pequenos.

write.index_bootstrap.tasks

Paralelismo da implantação

Paralelismo dos operadores de bootstrap de índice. Tem efeito apenas quando index.bootstrap.enabled é true. Aumentar este valor melhora o throughput do bootstrap, mas pode bloquear checkpoints durante o processo. Aumente a tolerância a falhas de checkpoint se necessário.

read.tasks

4

Paralelismo dos operadores de leitura em streaming e batch.

compaction.tasks

4

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

compaction.schedule.enabled

true

Define se os planos de compactação são gerados conforme agendamento. Mantenha como true mesmo com a compactação assíncrona desativada; assim, a compactação offline poderá executar os planos agendados.

compaction.async.enabled

true

Define se a compactação é executada de forma assíncrona. Defina como false para desativar a compactação online mantendo a geração de planos ativa.

compaction.tasks

4

Paralelismo das tarefas de compactação.

compaction.trigger.strategy

num_commits

Estratégia usada para acionar a compactação. Valores válidos: num_commits, time_elapsed, num_and_time, num_or_time.

compaction.delta_commits

5

Número de commits necessários para acionar a compactação. Usado com num_commits, num_and_time ou num_or_time.

compaction.delta_seconds

3600

Intervalo em segundos entre os acionamentos de compactação. Usado com time_elapsed, num_and_time ou num_or_time.

compaction.max_memory

100 MB

Memória máxima para o mapa hash usado durante a compactação e deduplicação. Aumente para 1 GB se os recursos permitirem.

compaction.target_io

500 GB

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

hoodie.parquet.max.file.size

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.

hoodie.parquet.small.file.limit

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.

hoodie.copyonwrite.record.size.estimate

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

hadoop.${option key}

Itens de configuração do Hadoop, especificados com o prefixo hadoop.. Suportado no Hudi 0.12.0 e posteriores. Use instruções DDL para especificar configurações do Hadoop por job em cenários entre clusters. Vários itens podem ser especificados simultaneamente.

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_insert ignora 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

write.operation

upsert

Tipo de gravação. Defina como bulk_insert para gravações em lote.

write.tasks

Paralelismo da implantação

Paralelismo para tarefas bulk_insert. O número final de arquivos de saída é maior ou igual a este valor (os dados passam para um novo arquivo quando o limite de 120 MB do Parquet é atingido).

write.bulk_insert.shuffle_input

true

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.

write.bulk_insert.sort_input

true

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.

write.sort.memory

128

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

changelog.enabled

false

Defina como true para reter todos os eventos de alteração. Quando false, apenas o registro final mesclado é garantido; alterações intermediárias podem ser combinadas.

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

write.insert.cluster

false

Defina como true para mesclar arquivos pequenos durante as gravações. Cada operação INSERT mescla arquivos pequenos existentes, mas não realiza deduplicação, e o throughput de gravação diminui.

Clustering assíncrono (Hudi 0.12.0 e posteriores)

Parâmetro

Padrão

Descrição

clustering.schedule.enabled

false

Defina como true para agendar periodicamente um plano de clustering.

clustering.delta_commits

4

Número de commits necessários para gerar um plano de clustering. Tem efeito apenas quando clustering.schedule.enabled é true.

clustering.async.enabled

false

Defina como true para executar o plano de clustering de forma assíncrona em intervalos regulares.

clustering.tasks

4

Paralelismo das tarefas de clustering.

clustering.plan.strategy.target.file.max.bytes

1 GiB (1.073.741.824 bytes)

Tamanho máximo alvo do arquivo para saída do clustering.

clustering.plan.strategy.small.file.limit

600

Arquivos menores que este limite (em bytes) são elegíveis para clustering.

clustering.plan.strategy.sort.columns

Colunas usadas para ordenar dados durante o clustering.

Estratégias de plano de clustering

Parâmetro

Padrão

Descrição

clustering.plan.partition.filter.mode

NONE

Modo de filtro de partição. Valores válidos: NONE (todas as partições), RECENT_DAYS (partições dos últimos N dias), SELECTED_PARTITIONS (partições específicas).

clustering.plan.strategy.daybased.lookback.partitions

2

Número de dias recentes para selecionar partições para clustering. Tem efeito apenas quando filter.mode é RECENT_DAYS.

clustering.plan.strategy.cluster.begin.partition

Partição inicial para filtragem por intervalo. Tem efeito apenas quando filter.mode é SELECTED_PARTITIONS.

clustering.plan.strategy.cluster.end.partition

Partição final para filtragem por intervalo. Tem efeito apenas quando filter.mode é SELECTED_PARTITIONS.

clustering.plan.strategy.partition.regex.pattern

Expressão regular para seleção de partições.

clustering.plan.strategy.partition.selected

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

Quando index.type está definido como BUCKET , definir index.global.enabled=true não tem efeito, pois o índice bucket não suporta deduplicação entre partições.

Parâmetro

Padrão

Descrição

index.type

FLINK_STATE

Tipo de índice. Valores válidos: FLINK_STATE, BUCKET.

hoodie.bucket.index.hash.field

Chave primária

Campo de chave hash para índice bucket. Pode ser um subconjunto da chave primária.

hoodie.bucket.index.num.buckets

4

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 read.streaming.enabled=true e opcionalmente read.start-commit

Leitura batch incremental

Defina tanto read.start-commit quanto read.end-commit; o intervalo é fechado (inclusivo em ambas as extremidades)

Time travel

Defina apenas read.end-commit; lê um snapshot naquele commit específico

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

read.streaming.enabled

false

Defina como true para ativar a leitura streaming.

read.start-commit

(vazio)

Offset inicial. Formato: yyyyMMddHHmmss para um horário específico ou earliest para ler desde o início. Deixe em branco para começar a partir do commit mais recente.

clean.retain_commits

30

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.

Importante

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

read.start-commit

Commit mais recente

Início do intervalo de leitura, no formato yyyyMMddHHmmss.

read.end-commit

Commit mais recente

Fim do intervalo de leitura, no formato yyyyMMddHHmmss. O intervalo é fechado (ambos os pontos finais inclusivos).

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

Importante

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>
Importante

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");
  }
}

Perguntas frequentes