Migra datos de captura de datos modificados (CDC) de Delta Lake a Lakehouse con Dataflow

Mantener sincronizadas las tablas de Apache Iceberg de Lakehouse de destino con las tablas de origen de Delta Lake 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 la creación de canalizaciones personalizadas para procesar los registros de transacciones introduce una sobrecarga operativa compleja.

Con las canalizaciones de captura de datos modificados (CDC) por lotes programadas de Dataflow, puedes sincronizar continuamente tus tablas de origen de Delta Lake con las tablas de Apache Iceberg de Lakehouse. En lugar de volver a cargar tablas completas, los trabajos por lotes programados extraen actualizaciones incrementales directamente de los registros de transacciones de Delta Lake y las aplican a tus tablas de destino en un programa configurable.

La canalización de CDC por lotes programada proporciona las siguientes capacidades:

  • Sincronización incremental rentable: Procesa periódicamente solo los datos modificados (operaciones INSERT, UPDATE y DELETE), lo que reduce la sobrecarga de procesamiento y la latencia de los datos.
  • Asignación automática de cambios: Analiza los registros del feed de datos de cambios de Delta Lake y asigna las operaciones de cambio a la tabla de Iceberg de Lakehouse de destino.
  • Programación flexible de trabajos: Ejecuta trabajos de Dataflow en modo por lotes según una programación (por ejemplo, intervalos de 15 minutos) que se adapte a las necesidades de tu empresa.

Flujo de trabajo de migración continua

Para configurar la migración continua de extremo a extremo de Delta Lake a Lakehouse, completa la siguiente secuencia de operaciones:

  1. Habilita la propiedad delta.enableChangeDataFeed (true) en la tabla de entrada de Delta Lake.
  2. Realiza una migración inicial de la tabla a un ID de confirmación o una marca de tiempo después de habilitar la propiedad delta.enableChangeDataFeed. Para obtener instrucciones, consulta Importa tablas de Delta Lake a Lakehouse con Dataflow.
  3. Ejecuta trabajos de CDC por lotes programados para actualizaciones de datos posteriores con una frecuencia (ID de confirmación o rango de marcas de tiempo) elegida para los requisitos de tu carga de trabajo.

Antes de comenzar

Para configurar la migración programada de CDC por lotes, asegúrate de tener lo siguiente:

  1. 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.

    Habilitar las API

  2. 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.

  3. Una tabla de Delta Lake válida almacenada en un bucket de Cloud Storage. El directorio de la tabla debe contener tus archivos de datos y el directorio del registro de transacciones _delta_log/. Además, la tabla debe cumplir con los siguientes requisitos:

    • La propiedad delta.enableChangeDataFeed debe estar habilitada (true) en la tabla de Delta Lake.
    • El ID de confirmación o la marca de tiempo de inicio configurados para la canalización de CDC deben ser posteriores a la habilitación de la propiedad delta.enableChangeDataFeed.
  4. 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 programada de CDC por lotes de Delta Lake 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 si proporcionas 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 origen de Delta Lake, debes actualizar manualmente el esquema de la tabla de destino de Lakehouse para que coincida antes de ejecutar el trabajo de sincronización de CDC.
  • Modo de ejecución por lotes: La CDC de Delta Lake opera en modo por lotes. La canalización extrae los registros de transacciones una vez por ejecución del trabajo entre los límites de confirmación de inicio y fin especificados (ID de confirmación o marca de tiempo). Si no se especifica un límite de confirmación final (end_version o end_timestamp), la canalización lee hasta la versión de confirmación más reciente.
  • Almacenamiento de origen compatible: Los datos de origen deben ser una tabla de Delta Lake válida almacenada en Cloud Storage (gs://). Amazon S3 no es compatible con las fuentes de tablas de Delta Lake.
  • Tablas basadas en catálogos: No se admiten las tablas de Delta Lake basadas en catálogos (por ejemplo, Unity Catalog).
  • Versiones no compatibles de Delta Lake: No se admite la versión 1.2.1 o anterior de Delta Lake porque el feed de datos de cambios (CDF) no es compatible con estas versiones.
  • Automatic Change Data Feed: No se admite Automatic Change Data Feed porque no escribe archivos de registro de cambios de CDF.

Crea una canalización de CDC de forma programática

Para compilar y ejecutar de forma programática una canalización de migración por lotes de CDC de Delta Lake, usa la transformación Managed.read en tu canalización de Apache Beam.

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 por lotes desde Delta Lake 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.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();
  }
}

Reemplaza lo siguiente:

  • BUCKET_NAME: Es el nombre del bucket de Cloud Storage que contiene la tabla de Delta Lake.
  • TABLE_NAME: Es el nombre del directorio de la tabla de origen de Delta Lake.
  • START_COMMIT_ID: Número de versión de la confirmación inicial para leer eventos de CDC de Delta Lake. Se debe proporcionar start_version o start_timestamp. El ID o la marca de tiempo de la confirmación de inicio deben ser posteriores a la habilitación de la propiedad delta.enableChangeDataFeed en la tabla de origen.
  • END_COMMIT_ID: Es el número de versión de la confirmación final para leer los eventos de CDC de Delta Lake. Como práctica recomendada, especifica un límite final (end_version o end_timestamp) que coincida con el tipo de límite que se usa para el inicio (start_version o start_timestamp). Si se omite un límite final, la canalización lee hasta la versión de confirmación más reciente.
  • 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.
  • PRIMARY_KEY_COLUMN: Es el nombre de la columna de clave primaria que se usa para identificar las filas de las actualizaciones y los borrados de CDC.

Examina el resultado del trabajo

Verifica que los datos de CDC se hayan combinado correctamente en tu tabla de Lakehouse:

  1. En la consola de Google Cloud , ve a la página de BigQuery Studio.

    Ir a BigQuery

  2. 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.
  3. Haz clic en Ejecutar y verifica los resultados.

¿Qué sigue?