Mantener sincronizadas las tablas de Apache Iceberg de Lakehouse de destino con las tablas de origen de Apache Iceberg que se actualizan con frecuencia puede ser difícil. Realizar recargas completas de tablas para capturar los cambios en los conjuntos de datos aumenta los costos de procesamiento y la latencia de procesamiento, y crear canalizaciones personalizadas para procesar los registros de captura de datos modificados (CDC) introduce una sobrecarga operativa compleja.
Con las canalizaciones de CDC por lotes o de transmisión de Dataflow, puedes sincronizar de forma continua o periódica tus tablas de origen de Apache Iceberg con las tablas de Apache Iceberg de Lakehouse. En lugar de volver a cargar tablas completas, las canalizaciones de CDC extraen actualizaciones incrementales directamente de los registros de cambios de Apache Iceberg de tu fuente y las aplican a tus tablas de destino.
La canalización de CDC proporciona las siguientes capacidades:
- Sincronización incremental rentable: Solo procesa los datos modificados (operaciones
INSERT,UPDATEyDELETE), lo que reduce la sobrecarga de procesamiento y la latencia de los datos. - Modos de ejecución flexibles: Admite modos por lotes y de transmisión para adaptarse a tus requisitos comerciales.
La canalización de CDC admite dos modos de ejecución:
- Modo de CDC por lotes: Lee periódicamente las actualizaciones incrementales de tu tabla de Apache Iceberg de origen entre los límites de instantáneas de inicio y finalización especificados.
- Modo de CDC de transmisión: Sondea continuamente tu tabla de Apache Iceberg de origen para detectar nuevas confirmaciones de cambios con una frecuencia configurable (de forma predeterminada, 1 minuto, con un mínimo de 1 segundo) y las aplica a tu tabla de Lakehouse casi en tiempo real.
Antes de comenzar
Para configurar la migración de CDC de Apache Iceberg a Lakehouse, asegúrate de tener lo siguiente:
Habilita las APIs de Dataflow, BigQuery y Lakehouse si alguna aún no está habilitada.
Roles necesarios para habilitar las APIs
Para habilitar APIs, necesitas el permiso
serviceusage.services.enable. Si creaste el proyecto, es probable que ya tengas este permiso a través del rol de propietario (roles/owner). De lo contrario, puedes obtener este permiso a través del rol de administrador de Service Usage (roles/serviceusage.serviceUsageAdmin). Obtén más información para otorgar roles.Para obtener los permisos que necesitas para crear los recursos, pídele a tu administrador que te otorgue los roles de Identity and Access Management (IAM) necesarios en tu proyecto.
Una tabla de Apache Iceberg de origen válida almacenada en Cloud Storage o registrada en un metastore o catálogo de REST de Iceberg.
Un catálogo de Iceberg de Lakehouse existente para recibir los datos sincronizados.
Si el espacio de nombres de destino no existe, la canalización de CDC lo crea automáticamente. Si la tabla de Lakehouse de destino no existe, la canalización de CDC la crea automáticamente si proporcionas una lista válida de columnas de igualdad.
Asistencia y limitaciones
La migración de CDC de Apache Iceberg a Lakehouse tiene las siguientes consideraciones:
- Requisito del SDK: Se requiere la versión 2.77.0 o posterior del SDK de Apache Beam.
- Creación de la tabla de destino: Se admite la creación de la tabla de Lakehouse de destino. Si la tabla de destino no existe, la canalización de CDC la crea automáticamente, siempre que proporciones una lista válida de columnas de igualdad (
equality_columns). Solo se requiere un catálogo de Lakehouse existente; la canalización de CDC crea automáticamente el espacio de nombres si no existe. - Evolución del esquema: No se admiten los cambios automáticos en el esquema. Si cambia el esquema de la tabla de Apache Iceberg de origen, debes actualizar manualmente el esquema de la tabla de Lakehouse de destino para que coincida antes de ejecutar el trabajo de sincronización de CDC.
- Modos de ejecución:
- CDC por lotes: Extrae los cambios del registro de transacciones una vez por ejecución del trabajo entre los límites de instantánea de inicio y finalización especificados (ID de instantánea o marca de tiempo).
- CDC de transmisión: Consulta continuamente la tabla de Apache Iceberg de origen para obtener registros de confirmación nuevos. La frecuencia de sondeo se puede configurar (el valor predeterminado es de 1 minuto y el mínimo es de 1 segundo).
Crea una canalización de CDC con Java
Para compilar y ejecutar de forma programática una canalización de migración por lotes o de transmisión de CDC de Apache Iceberg con el SDK de Apache Beam Java, usa la transformación Managed.read con Managed.ICEBERG_CDC.
En los siguientes ejemplos de Java, se muestra cómo configurar y ejecutar canalizaciones de CDC de Apache Iceberg para cargas de trabajo por lotes y de transmisión.
Canalización de CDC por lotes
En el siguiente ejemplo de Java, se configura una canalización por lotes que lee eventos de CDC de Apache Iceberg dentro de los límites de instantáneas de inicio y finalización especificados:
Java
Agrega las siguientes dependencias a tu archivo 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>
En el siguiente fragmento de código en Java, se muestra cómo configurar y ejecutar una canalización de CDC por lotes desde una tabla de Iceberg hasta 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();
}
}
Reemplaza lo siguiente:
SOURCE_CATALOG_TYPE: Es el tipo del catálogo de Iceberg de origen (por ejemplo,hadoopohive).SOURCE_WAREHOUSE_LOCATION: Es la ubicación del almacén del catálogo de Iceberg de origen (por ejemplo,s3://source-warehouseogs://source-warehouse).SOURCE_NAMESPACE: Es el espacio de nombres de la tabla Iceberg de origen.SOURCE_TABLE: Es el nombre de la tabla de Iceberg de origen.START_SNAPSHOT_ID: Es el ID de la instantánea inicial para la lectura de CDC por lotes.END_SNAPSHOT_ID: Es el ID de la instantánea final para la lectura de CDC por lotes.WAREHOUSE_BUCKET: Es el nombre del bucket de Cloud Storage que se usa como almacén del catálogo de Lakehouse.PROJECT_ID: Es el ID del proyecto de Google Cloud .TARGET_NAMESPACE: Es el espacio de nombres de la tabla de Lakehouse de destino.TARGET_TABLE: Es el nombre de la tabla de Lakehouse de destino.
Canalización de CDC de transmisión
En el siguiente ejemplo de Java, se configura una canalización de transmisión que sondea continuamente la tabla de Apache Iceberg de origen para detectar nuevos cambios de CDC:
Java
Agrega las siguientes dependencias a tu archivo 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>
En el siguiente fragmento de Java, se muestra cómo configurar y ejecutar una canalización de CDC de transmisión desde una tabla de Iceberg hasta 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();
}
}
Reemplaza lo siguiente:
SOURCE_CATALOG_TYPE: Es el tipo del catálogo de Iceberg de origen (por ejemplo,hadoopohive).SOURCE_WAREHOUSE_LOCATION: Es la ubicación del almacén del catálogo de Iceberg de origen (por ejemplo,s3://source-warehouseogs://source-warehouse).SOURCE_NAMESPACE: Es el espacio de nombres de la tabla Iceberg de origen.SOURCE_TABLE: Es el nombre de la tabla de Iceberg de origen.POLL_INTERVAL_SECONDS: Es la frecuencia de sondeo continuo en segundos (el valor predeterminado es60y el mínimo es1).WAREHOUSE_BUCKET: Es el nombre del bucket de Cloud Storage que se usa como almacén del catálogo de Lakehouse.PROJECT_ID: Es el ID del proyecto de Google Cloud .TARGET_NAMESPACE: Es el espacio de nombres de la tabla de Lakehouse de destino.TARGET_TABLE: Es el nombre de la tabla de Lakehouse de destino.
Examina el resultado del trabajo
Verifica que los datos de CDC se hayan combinado correctamente en tu tabla de Lakehouse:
En la consola de Google Cloud , ve a la página de BigQuery Studio.
En el editor de consultas, ejecuta una consulta en SQL para verificar los datos sincronizados:
SELECT * FROM `PROJECT_ID`.`CATALOG`.`NAMESPACE`.`TABLE_NAME` LIMIT 10;Reemplaza lo siguiente:
PROJECT_ID: Es el ID del proyecto de Google Cloud .CATALOG: Es el nombre de tu catálogo de Lakehouse.NAMESPACE: Es el espacio de nombres de tu tabla de Lakehouse.TABLE_NAME: Es el nombre de la tabla de Lakehouse de destino.
Haz clic en Ejecutar y verifica los resultados.
¿Qué sigue?
- Obtén más información para importar tablas de Iceberg a Lakehouse.
- Obtén más información sobre la ingesta de captura de datos modificados en Lakehouse.
- Obtén más información sobre la E/S administrada en Dataflow.