Manter as tabelas de destino do Lakehouse Apache Iceberg sincronizadas com tabelas de origem do Delta Lake atualizadas com frequência pode ser difícil. A execução de recarregamentos completos de tabelas para capturar mudanças no conjunto de dados aumenta os custos de computação e a latência de processamento. Além disso, a criação de pipelines personalizados para processar registros de transações introduz uma sobrecarga operacional complexa.
Ao usar pipelines de captura de dados alterados (CDC) em lote programados do Dataflow, é possível sincronizar continuamente as tabelas de origem do Delta Lake com as tabelas do Lakehouse Apache Iceberg. Em vez de recarregar tabelas completas, os jobs em lote programados extraem atualizações incrementais diretamente dos registros de transações do Delta Lake e as aplicam às tabelas de destino em uma programação configurável.
O pipeline de CDC em lote programado oferece os seguintes recursos:
- Sincronização incremental econômica: processa periodicamente apenas os dados alterados (operações
INSERT,UPDATEeDELETE), reduzindo a sobrecarga de computação e a latência de dados. - Mapeamento automático de mudanças: analisa os registros do Change Data Feed do Delta Lake e mapeia as operações de mudança para a tabela de destino do Lakehouse Iceberg.
- Programação flexível de jobs: executa jobs do Dataflow no modo em lote em uma programação (por exemplo, intervalos de 15 minutos) que atenda às suas necessidades comerciais.
Fluxo de trabalho de migração contínua
Para configurar a migração contínua de ponta a ponta do Delta Lake para o Lakehouse, conclua a seguinte sequência de operações:
- Ative a propriedade
delta.enableChangeDataFeed(true) na tabela de entrada do Delta Lake. - Faça uma migração inicial de tabela para um ID de commit ou carimbo de data/hora depois de
ativar a propriedade
delta.enableChangeDataFeed. Para instruções, consulte Importar tabelas do Delta Lake para o Lakehouse usando o Dataflow. - Execute jobs de CDC em lote programados para atualizações de dados subsequentes com uma frequência (ID de commit ou intervalo de carimbos de data/hora) escolhida para os requisitos da sua carga de trabalho.
Antes de começar
Para configurar a migração de CDC em lote programada, verifique se você tem o seguinte:
Ative as APIs Dataflow, BigQuery e Lakehouse, se alguma delas ainda não estiver ativada.
Funções necessárias para ativar APIs
Para ativar APIs, você precisa da permissão
serviceusage.services.enable. Se você criou o projeto, provavelmente já tem essa permissão pelo papel de proprietário (roles/owner). Caso contrário, é possível receber essa permissão pelo papel de administrador do Service Usage (roles/serviceusage.serviceUsageAdmin). Saiba como conceder papéis.Para receber as permissões necessárias para criar os recursos, peça ao administrador para conceder a você os papéis necessários do Identity and Access Management (IAM) no projeto.
Uma tabela válida do Delta Lake armazenada em um bucket do Cloud Storage. O diretório de tabela precisa conter seus arquivos de dados e o diretório de registro de transações
_delta_log/. Além disso, a tabela precisa atender aos seguintes requisitos:- A propriedade
delta.enableChangeDataFeedprecisa estar ativada (true) na tabela do Delta Lake. - O ID do commit inicial ou o carimbo de data/hora configurado para o pipeline de CDC precisa ser posterior à ativação da propriedade
delta.enableChangeDataFeed.
- A propriedade
Um catálogo do Lakehouse Iceberg para receber os dados sincronizados. Se o namespace de destino não existir, o pipeline de CDC vai criá-lo automaticamente. Se a tabela de destino do Lakehouse não existir, o pipeline de CDC a criará automaticamente se você fornecer uma lista válida de colunas de igualdade.
Suporte e limitações
A migração programada de CDC em lote do Delta Lake para o Lakehouse tem as seguintes considerações:
- Requisito do SDK:requer o SDK do Apache Beam versão 2.77.0 ou mais recente.
- Criação da tabela de destino:é possível criar a tabela de destino do Lakehouse. Se a tabela de destino não existir, o pipeline de CDC a criará automaticamente se você fornecer uma lista válida de colunas de igualdade (
equality_columns). Apenas um catálogo do Lakehouse é necessário. O pipeline de CDC cria automaticamente o namespace se ele não existir. - Evolução do esquema:não há suporte para mudanças automáticas de esquema. Se o esquema da tabela de origem do Delta Lake mudar, atualize manualmente o esquema da tabela de destino do Lakehouse para corresponder antes de executar o job de sincronização de CDC.
- Modo de execução em lote:o CDC do Delta Lake opera no modo em lote.
O pipeline extrai os registros de transações uma vez por execução de job entre os limites de commit inicial e final especificados (ID ou carimbo de data/hora do commit). Se um limite de confirmação final (
end_versionouend_timestamp) não for especificado, o pipeline vai ler até a versão de confirmação mais recente. - Armazenamento de origem compatível:os dados de origem precisam ser uma tabela Delta Lake válida armazenada no Cloud Storage (
gs://). O Amazon S3 não é compatível com fontes de tabela Delta Lake. - Tabelas baseadas em catálogo:as tabelas Delta Lake baseadas em catálogo (por exemplo, Unity Catalog) não são compatíveis.
- Versões sem suporte do Delta Lake:a versão 1.2.1 ou anterior do Delta Lake não é compatível porque o feed de dados de mudança (CDF) não é compatível com essas versões.
- Feed de dados de mudança automático:não há suporte para o feed de dados de mudança automático porque ele não grava arquivos de registro de alterações do CDF.
Criar um pipeline de CDC de maneira programática
Para criar e executar programaticamente um pipeline de migração em lote de CDC do Delta Lake, use a transformação Managed.read no pipeline do Apache Beam.
Java
Adicione as seguintes dependências ao arquivo pom.xml:
<dependency>
<groupId>org.apache.beam</groupId>
<artifactId>beam-sdks-java-managed</artifactId>
<version>${beam.version}</version>
</dependency>
<dependency>
<groupId>org.apache.beam</groupId>
<artifactId>beam-sdks-java-io-iceberg</artifactId>
<version>${beam.version}</version>
</dependency>
O snippet Java a seguir demonstra como configurar e executar um pipeline de CDC em lote do Delta Lake para o Lakehouse:
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.managed.Managed;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
public class DeltaLakeCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// 1. Configure Delta Lake CDC reader
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "gs://BUCKET_NAME/TABLE_NAME");
// Starting commit bound (either start_version or start_timestamp must be provided):
readConfig.put("start_version", START_COMMIT_ID);
// Optional ending commit bound (end_version or end_timestamp, matching the start bound type):
// readConfig.put("end_version", END_COMMIT_ID);
readConfig.put("include_metadata_columns", Arrays.asList("_commit_version", "_change_type"));
// Read CDC events from Delta Lake
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromDeltaCDC",
Managed.read(Managed.DELTA_LAKE_CDC).withConfig(readConfig)).getSinglePCollection();
// 2. Configure BigLake REST Catalog properties for Lakehouse Iceberg
Map<String, String> catalogProps = new HashMap<>();
catalogProps.put("type", "rest");
catalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
catalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
catalogProps.put("header.x-goog-user-project", "PROJECT_ID");
catalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
catalogProps.put("rest.auth.type", "org.apache.iceberg.gcp.auth.GoogleAuthManager");
// 3. Map Delta Lake change types to Iceberg CDC change types
Map<String, String> changeTypeMap = new HashMap<>();
changeTypeMap.put("insert", "INSERT");
changeTypeMap.put("delete", "DELETE");
changeTypeMap.put("update_preimage", "UPDATE_BEFORE");
changeTypeMap.put("update_postimage", "UPDATE_AFTER");
// 4. Configure Lakehouse Iceberg CDC writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_name", "lakehouse");
writeConfig.put("catalog_properties", catalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("sequence_number_column", "_commit_version");
writeConfig.put("equality_columns", Arrays.asList("PRIMARY_KEY_COLUMN"));
// change_type_column and change_type_map are only required for Portable Runner
writeConfig.put("change_type_column", "_change_type");
writeConfig.put("change_type_map", changeTypeMap);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
Substitua:
BUCKET_NAME: o nome do bucket do Cloud Storage que contém a tabela do Delta Lake.TABLE_NAME: o nome do diretório da tabela de origem do Delta Lake.START_COMMIT_ID: número da versão inicial do commit para leitura de eventos de CDC do Delta Lake. É necessário fornecerstart_versionoustart_timestamp. O ID ou carimbo de data/hora do commit inicial precisa ser posterior à ativação da propriedadedelta.enableChangeDataFeedna tabela de origem.END_COMMIT_ID: número da versão do commit final para leitura de eventos de CDC do Delta Lake. Como prática recomendada, especifique um limite final (end_versionouend_timestamp) que corresponda ao tipo de limite usado para início (start_versionoustart_timestamp). Se um limite final for omitido, o pipeline vai ler até a versão de commit mais recente.WAREHOUSE_BUCKET: o nome do bucket do Cloud Storage usado como o data warehouse do catálogo do Lakehouse.PROJECT_ID: o ID do projeto Google Cloud .TARGET_NAMESPACE: o namespace da tabela de destino do Lakehouse.TARGET_TABLE: o nome da tabela de destino do Lakehouse.PRIMARY_KEY_COLUMN: o nome da coluna de chave primária usada para identificar linhas em atualizações e exclusões de CDC.
Examinar a saída do job
Verifique se os dados da CDC foram mesclados com sucesso na tabela do Lakehouse:
No console Google Cloud , acesse a página Studio do BigQuery.
No editor de consultas, execute uma consulta SQL para verificar os dados sincronizados:
SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;Substitua:
PROJECT_ID: o ID do projeto Google Cloud .CATALOG: o nome do catálogo do Lakehouse.NAMESPACE: o namespace da sua tabela do Lakehouse.TABLE_NAME: o nome da tabela de destino do Lakehouse.
Clique em Executar e verifique os resultados.
A seguir
- Saiba mais sobre como importar tabelas do Delta Lake para o Lakehouse.
- Aprenda mais sobre a ingestão de captura de dados alterados no Lakehouse.
- Saiba mais sobre a E/S gerenciada no Dataflow.