Manter as tabelas de destino do Lakehouse Apache Iceberg sincronizadas com tabelas de origem do Apache Iceberg 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 captura de dados alterados (CDC) introduz uma sobrecarga operacional complexa.
Ao usar pipelines de CDC em lote ou de streaming do Dataflow, é possível sincronizar continuamente ou periodicamente as tabelas de origem do Apache Iceberg com as tabelas do Apache Iceberg do Lakehouse. Em vez de recarregar tabelas completas, os pipelines de CDC extraem atualizações incrementais diretamente dos registros de alterações do Apache Iceberg de origem e as aplicam às tabelas de destino.
O pipeline de CDC oferece os seguintes recursos:
- Sincronização incremental econômica: processa apenas dados alterados (operações
INSERT,UPDATEeDELETE), reduzindo a sobrecarga de computação e a latência de dados. - Modos de execução flexíveis: compatível com os modos em lote e de streaming para atender aos requisitos da sua empresa.
O pipeline de CDC é compatível com dois modos de execução:
- Modo CDC em lote:lê periodicamente atualizações incrementais da tabela de origem do Apache Iceberg entre os limites de início e fim especificados do snapshot.
- Modo de CDC de streaming:pesquisa continuamente a tabela de origem do Apache Iceberg para novos commits de mudança em uma frequência configurável (o padrão é 1 minuto, mínimo de 1 segundo) e os aplica à tabela do Lakehouse quase em tempo real.
Antes de começar
Para configurar a migração de CDC do Apache Iceberg para o Lakehouse, 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 de origem válida do Apache Iceberg armazenada no Cloud Storage ou registrada em um catálogo REST ou metastore do Iceberg.
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 de CDC do Apache Iceberg 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, desde que você forneça uma lista válida de colunas de igualdade (
equality_columns). Apenas um catálogo do Lakehouse existente é 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 Apache Iceberg mudar, atualize manualmente o esquema da tabela de destino do Lakehouse para corresponder antes de executar o job de sincronização de CDC.
- Modos de execução:
- CDC em lote:extrai as mudanças do registro de transações uma vez por execução do job entre os limites de snapshot inicial e final especificados (ID ou carimbo de data/hora do snapshot).
- CDC de streaming:pesquisa continuamente a tabela de origem do Apache Iceberg em busca de novos registros de commit. A frequência de sondagem é configurável (o padrão é de 1 minuto, com um mínimo de 1 segundo).
Criar um pipeline de CDC usando Java
Para criar e executar de maneira programática um pipeline de migração em lote ou de streaming de CDC do Apache Iceberg usando o SDK do Apache Beam para Java, use a transformação Managed.read com Managed.ICEBERG_CDC.
Os exemplos de Java a seguir demonstram como configurar e executar pipelines de CDC do Apache Iceberg para cargas de trabalho em lote e de streaming.
Pipeline de CDC em lote
O exemplo em Java a seguir configura um pipeline em lote que lê eventos de CDC do Apache Iceberg dentro dos limites de início e fim especificados do snapshot:
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 de Java a seguir demonstra como configurar e executar um pipeline de CDC em lote de uma tabela do Iceberg 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergBatchCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// Configure source catalog properties
Map<String, String> sourceCatalogProps = new HashMap<>();
sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
// Add additional properties if required by your catalog type (for example, uri)
// Configure batch Iceberg CDC reader with snapshot bounds
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
readConfig.put("catalog_properties", sourceCatalogProps);
readConfig.put("start_snapshot_id", START_SNAPSHOT_ID);
readConfig.put("end_snapshot_id", END_SNAPSHOT_ID);
readConfig.put(
"include_metadata_columns",
Arrays.asList(
IcebergCdcMetadataColumns.CHANGE_TYPE,
IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));
// Read batch CDC events from source Iceberg table
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromIcebergCDC",
Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();
// Configure destination Lakehouse REST Catalog properties for Lakehouse Iceberg
Map<String, String> destinationCatalogProps = new HashMap<>();
destinationCatalogProps.put("type", "rest");
destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
destinationCatalogProps.put("rest.auth.type", "google");
// Configure Lakehouse Iceberg writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_properties", destinationCatalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
Substitua:
SOURCE_CATALOG_TYPE: o tipo do catálogo de origem do Iceberg (por exemplo,hadoopouhive).SOURCE_WAREHOUSE_LOCATION: o local do depósito do catálogo de origem do Iceberg (por exemplo,s3://source-warehouseougs://source-warehouse).SOURCE_NAMESPACE: o namespace da tabela Iceberg de origem.SOURCE_TABLE: o nome da tabela de origem do Iceberg.START_SNAPSHOT_ID: o ID do snapshot inicial para a leitura em lote da CDC.END_SNAPSHOT_ID: o ID do snapshot final para a leitura de CDC em lote.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.
Pipeline de CDC de streaming
O exemplo em Java a seguir configura um pipeline de streaming que pesquisa continuamente a tabela de origem do Apache Iceberg para novas mudanças de CDC:
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 de streaming de uma tabela do Iceberg 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.io.iceberg.cdc.IcebergCdcMetadataColumns;
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 IcebergStreamingCdcToLakehouse {
public static void main(String[] args) {
Pipeline pipeline = Pipeline.create(PipelineOptionsFactory.fromArgs(args).create());
// Configure source catalog properties
Map<String, String> sourceCatalogProps = new HashMap<>();
sourceCatalogProps.put("type", "SOURCE_CATALOG_TYPE");
sourceCatalogProps.put("warehouse", "SOURCE_WAREHOUSE_LOCATION");
// Add additional properties if required by your catalog type (for example, uri)
// Configure streaming Iceberg CDC reader with continuous polling
Map<String, Object> readConfig = new HashMap<>();
readConfig.put("table", "SOURCE_NAMESPACE.SOURCE_TABLE");
readConfig.put("catalog_properties", sourceCatalogProps);
readConfig.put("streaming", true);
readConfig.put("poll_interval_seconds", POLL_INTERVAL_SECONDS);
readConfig.put(
"include_metadata_columns",
Arrays.asList(
IcebergCdcMetadataColumns.CHANGE_TYPE,
IcebergCdcMetadataColumns.COMMIT_SNAPSHOT_SEQUENCE_NUMBER));
// Read streaming CDC events from source Iceberg table
PCollection<Row> cdcRows =
pipeline.apply(
"ReadFromIcebergCDC",
Managed.read(Managed.ICEBERG_CDC).withConfig(readConfig)).getSinglePCollection();
// Configure destination BigLake REST Catalog properties for Lakehouse Iceberg
Map<String, String> destinationCatalogProps = new HashMap<>();
destinationCatalogProps.put("type", "rest");
destinationCatalogProps.put("uri", "https://biglake.googleapis.com/iceberg/v1/restcatalog");
destinationCatalogProps.put("warehouse", "gs://WAREHOUSE_BUCKET");
destinationCatalogProps.put("header.x-goog-user-project", "PROJECT_ID");
destinationCatalogProps.put("io-impl", "org.apache.iceberg.gcp.gcs.GCSFileIO");
destinationCatalogProps.put("rest.auth.type", "google");
// Configure Lakehouse Iceberg writer
Map<String, Object> writeConfig = new HashMap<>();
writeConfig.put("table", "TARGET_NAMESPACE.TARGET_TABLE");
writeConfig.put("catalog_properties", destinationCatalogProps);
writeConfig.put("mode", "merge-on-read");
writeConfig.put("change_type_column", IcebergCdcMetadataColumns.CHANGE_TYPE);
// Write CDC rows to target Lakehouse Iceberg table
cdcRows.apply(
"WriteToLakehouse",
Managed.write(Managed.ICEBERG).withConfig(writeConfig));
pipeline.run();
}
}
Substitua:
SOURCE_CATALOG_TYPE: o tipo do catálogo de origem do Iceberg (por exemplo,hadoopouhive).SOURCE_WAREHOUSE_LOCATION: o local do depósito do catálogo de origem do Iceberg (por exemplo,s3://source-warehouseougs://source-warehouse).SOURCE_NAMESPACE: o namespace da tabela Iceberg de origem.SOURCE_TABLE: o nome da tabela de origem do Iceberg.POLL_INTERVAL_SECONDS: a frequência de sondagem contínua em segundos (padrão60, mínimo1).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.
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 Iceberg 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.