Temas avanzados de AUTO CDC

Además de las API básicas AUTO CDC y AUTO CDC FROM SNAPSHOT, puede ejecutar DML en tablas de destino, leer fuentes de datos de cambios de destinos de CDC, supervisar métricas de procesamiento, aplicar actualizaciones parciales y realizar un seguimiento de los cambios con almacenamiento bitemporal. Para obtener una introducción a las AUTO CDC APIs, consulte Las API de AUTO CDC: Simplifica la captura de datos cambiantes con canalizaciones.

Agregar, cambiar o eliminar datos en una tabla de streaming de destino

Si la canalización publica tablas en el catálogo de Unity, use instrucciones de lenguaje de manipulación de datos (DML), incluyendo las instrucciones insert, update, delete y merge, para modificar las tablas de streaming de destino creadas por instrucciones AUTO CDC ... INTO.

Nota:

  • No se admiten instrucciones DML que modifican el esquema de tabla de una tabla de streaming. Asegúrese de que las instrucciones DML no intenten evolucionar el esquema de la tabla.
  • Las instrucciones DML que actualizan una tabla de streaming solo se pueden ejecutar en un clúster compartido de Unity Catalog o en un almacenamiento de SQL mediante Databricks Runtime 13.3 LTS y versiones posteriores.
  • Dado que el streaming requiere orígenes de datos de solo anexión, si el procesamiento requiere streaming desde una tabla de streaming de origen con cambios (por ejemplo, por instrucciones de DML), establezca la marca skipChangeCommits al leer la tabla de streaming de origen. Cuando se establezca skipChangeCommits, se omitirán las transacciones que eliminen o modifiquen registros de la tabla de origen. Si el procesamiento no requiere una tabla de streaming, puede usar una vista materializada (que no tiene la restricción de solo anexión) como tabla de destino.

Dado que la canalización usa una columna especificada SEQUENCE BY y propaga los valores de secuenciación adecuados a las __START_AT columnas y __END_AT de la tabla de destino (para SCD Type 2), debe asegurarse de que las instrucciones DML usan valores válidos para estas columnas para mantener el orden adecuado de los registros. Consulte Funcionamiento de AUTO CDC.

Para obtener más información sobre el uso de instrucciones DML con tablas de streaming, vea Agregar, cambiar o eliminar datos en una tabla de streaming.

En el ejemplo siguiente se inserta un registro activo con una secuencia de inicio de 5:

INSERT INTO my_streaming_table (id, name, __START_AT, __END_AT) VALUES (123, 'John Doe', 5, NULL);

Sugerencia

Si necesita cambiar el nombre de las columnas __START_AT y __END_AT en su tabla de destino SCD Tipo 2 (por ejemplo, para que coincida con los requisitos del esquema subsiguiente), cree una vista de la tabla de destino:

CREATE VIEW my_employees_view AS
SELECT
  *,
  __START_AT AS valid_from,
  __END_AT AS valid_to
FROM my_scd2_target_table;

Leer una fuente de datos modificados de una tabla de destino AUTO CDC

En Databricks Runtime 15.2 y versiones posteriores, puede leer una fuente de datos de cambios de una tabla de streaming que sea el destino de consultas AUTO CDC o AUTO CDC FROM SNAPSHOT de la misma manera que lee una fuente de datos de cambios de otras tablas Delta. Se requiere lo siguiente para leer la fuente de datos de cambios de una tabla de streaming de destino:

  • La tabla de streaming de destino debe publicarse en el catálogo de Unity. Consulta Usar Unity Catalog con canalizacións.
  • Para leer el flujo de cambios de datos desde la tabla de transmisión de destino, debe usar Databricks Runtime 15.2 o superior. Para leer el flujo de datos de cambios en una canalización diferente, la canalización debe configurarse para usar Databricks Runtime 15.2 o superior.

La lectura del flujo de datos modificados desde una tabla de streaming de destino creada en una canalización de Lakeflow se realiza de la misma forma que la lectura de un flujo de datos modificados desde otras tablas de Delta. Para obtener más información sobre cómo usar la funcionalidad de fuente de datos de cambios de Delta, incluidos ejemplos en Python y SQL, consulte Uso de la fuente de datos de cambios en Azure Databricks.

Nota:

El registro de fuente de distribución de datos modificados incluye metadatos que identifican el tipo de evento de cambio. Cuando se actualiza un registro en una tabla, los metadatos de los registros de cambios asociados suelen incluir valores de _change_type establecidos en eventos update_preimage y update_postimage.

Sin embargo, los _change_type valores son diferentes si se realizan actualizaciones en la tabla de streaming de destino que incluyen el cambio de valores de clave principal. Cuando los cambios incluyen actualizaciones de las claves principales, los _change_type campos de metadatos se establecen en eventos de insert y delete. Los cambios en las claves principales pueden ocurrir cuando se realizan actualizaciones manuales en uno de los campos clave con una instrucción UPDATE o MERGE, o, para las tablas de tipo SCD 2, cuando el campo __start_at cambia para reflejar un valor de secuencia de inicio anterior.

La consulta AUTO CDC determina los valores de la clave primaria, que difieren según se trate de un procesamiento SCD de tipo 1 o de tipo 2:

Tipo SCD Clave principal
SCD de tipo 1 y la interfaz Python de los pipelines La clave principal es el valor del keys parámetro en la create_auto_cdc_flow() función . Para la interfaz SQL, la clave principal es las columnas definidas por la KEYS cláusula en la AUTO CDC ... INTO instrucción .
SCD de tipo 2 La clave principal es el keys parámetro o KEYS cláusula más el valor devuelto de la coalesce(__START_AT, __END_AT) operación, donde __START_AT y __END_AT son las columnas correspondientes de la tabla de streaming de destino. Esto usa __START_AT cuando está disponible y __END_AT cuando __START_AT es null (por ejemplo, el registro inicial).

Leer un flujo de datos de cambios de una vista materializada

Importante

Esta característica se encuentra en su versión beta.

Puedes leer un feed de datos de cambios desde una vista materializada creada en una tubería de Lakeflow o en Databricks SQL. Úsalo para replicar cambios materializados en vistas a destinos fuera de Azure Databricks, o para mantener un historial de cambios materializados en vistas para auditoría e informes.

Las vistas materializadas usan una fuente automática de datos de cambios, por lo que no es necesario habilitar la fuente de datos de cambios en sí. En su lugar, habilitas la fuente de datos de cambios en cada vista materializada en la que necesites habilitarla si cumples los siguientes requisitos. Consulte Fuente de distribución automática de datos de cambios.

  • Para leer el feed de datos de cambios, debes usar Databricks Runtime 18 LTS o superior, en computación clásica, computación serverless o Databricks SQL.

  • La vista materializada, la canalización que la crea o la canalización que la lee debe usar el canal PREVIEW.

  • La vista materializada debe tener activado el seguimiento de filas. Las vistas materializadas en cómputo sin servidor tienen habilitado el seguimiento de filas por defecto. Consulte Seguimiento de filas en Azure Databricks. Para comprobar si el seguimiento de filas está habilitado en una vista materializada, ejecuta:

    SHOW TBLPROPERTIES my_mv ('delta.enableRowTracking');
    
  • Para leer el flujo de datos de cambios de una vista materializada, activa la marca de metadatos externos en la canalización o en la vista materializada. Para las instrucciones, consulte Cómo habilitar el acceso a un conjunto de datos.

Lees el feed de datos de cambio desde una vista materializada igual que desde otras tablas Delta, usando la table_changes() función, una lectura en streaming o la readChangeFeed opción. Para la sintaxis y los ejemplos en SQL y Python, consulta Usar la fuente de cambios de datos en Azure Databricks.

Puedes leer un flujo de datos de cambio de vista materializado desde dentro de una vista materializada SQL de Databricks o una tabla de streaming:

CREATE OR REFRESH STREAMING TABLE sales
  AS SELECT * FROM STREAM my_mv WITH (readChangeFeed=true)

Limitaciones

Además de las limitaciones de la alimentación automática de datos de cambio, se aplican las siguientes cosas cuando lees una fuente de datos de cambio desde una vista materializada:

  • El feed de datos de cambio incluye filas sin cambios cuando la vista materializada se reescribe completamente, y no consolida múltiples actualizaciones de la misma fila en un solo evento. Para filtrarlos, agrupe el flujo de datos de cambios por todas las columnas para encontrar inserciones y eliminaciones que tengan los mismos valores de fila.
  • Solo Azure Databricks puede consultar el feed de datos de cambios para obtener una vista materializada. Los clientes externos de Delta Lake e Iceberg no pueden.
  • Dentro de las canalizaciones de Lakeflow, solo puedes leer el flujo de datos de cambios de una vista materializada desde una canalización diferente, y esa canalización debe usar el canal PREVIEW. No se admite leer el flujo de datos de cambios de una vista materializada en la misma canalización que la crea.
  • No se puede crear un índice de búsqueda vectorial a partir de una vista materializada.

Obtener datos sobre los registros procesados por una consulta CDC en los pipelines

Nota:

Las siguientes métricas solo las capturan las AUTO CDC consultas y no las AUTO CDC FROM SNAPSHOT consultas.

Las siguientes métricas son capturadas por las consultas AUTO CDC.

  • num_upserted_rows: El número de filas de salida insertadas o actualizadas en el conjunto de datos durante una actualización.
  • num_deleted_rows: número de filas de salida existentes eliminadas del conjunto de datos durante una actualización.

La métrica num_output_rows, que es la salida de los flujos que no son CDC, no se captura para las consultas AUTO CDC.

Aplicar actualizaciones parciales

Cuando un origen envía solo las columnas que han cambiado, AUTO CDC debe distinguir entre una columna ausente de un registro de cambio, que debe dejar el valor de destino sin cambios y una columna que se establece nullexplícitamente en , que debe sobrescribir el valor de destino con null. De forma predeterminada, IGNORE NULL UPDATES trata cada null como un marcador "no actualizar", por lo que no puede aplicar un elemento explícito null. Para resolver esta ambigüedad, elija uno de los tres métodos siguientes:

Método Cuándo usarlo Comportamiento
IGNORE NULL UPDATES ON columnList Un conjunto pequeño y fijo de columnas debe omitir null los valores, mientras que todas las demás columnas aplican valores explícitos null . Las columnas enumeradas mantienen su valor de destino existente cuando el valor entrante es null. Todas las demás columnas aplican valores explícitos null .
IGNORE NULL UPDATES ON * EXCEPT (exceptColumnList) La mayoría de las columnas deberían omitir los valores null, y solo unas pocas deberían aplicar valores null explícitos. Las columnas enumeradas aplican valores explícitos null . Todas las demás columnas mantienen su valor de destino existente cuando el valor entrante es null.
COLUMNS TO UPDATE Cada registro de cambios actualiza un conjunto diferente de columnas, o bien el conjunto de columnas actualizables cambia con el tiempo. Una columna de origen asigna un nombre a las columnas que se van a actualizar para cada registro de cambio. Las columnas listadas se escriben a partir de la fuente, incluidos los valores explícitos null. Las columnas que no aparecen en la lista mantienen su valor de destino existente.

COLUMNS TO UPDATE no se puede combinar con IGNORE NULL UPDATESy no se admite para tablas bitemporales.

Como regla general, elija COLUMNS TO UPDATE cuándo el productor sabe qué columnas han cambiado en cada registro y puede llevar esa información en una columna de origen, como cuando varios productores escriben en el mismo origen o el conjunto de columnas actualizables crece con el tiempo. Elija IGNORE NULL UPDATES ON cuándo el propietario de la canalización conoce el conjunto fijo de columnas actualizables de antemano y prefiere controlarlas en el código de canalización.

En el ejemplo siguiente se usa una columna de origen denominada columnsToUpdate para controlar qué columnas actualiza cada registro de cambios, incluidas las columnas establecidas explícitamente en null:

Python

from pyspark import pipelines as dp

dp.create_streaming_table("target")

dp.create_auto_cdc_flow(
  target = "target",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "sequenceNum",
  stored_as_scd_type = 1,
  columns_to_update = "columnsToUpdate"
)

SQL

CREATE OR REFRESH STREAMING TABLE target;

CREATE FLOW apply_cdc AS AUTO CDC INTO
  target
FROM
  stream(cdc_source)
KEYS
  (id)
SEQUENCE BY
  sequenceNum
STORED AS
  SCD TYPE 1
COLUMNS TO UPDATE
  columnsToUpdate;

Para obtener la referencia completa de parámetros, consulte AUTO CDC INTO (pipelines) y create_auto_cdc_flow.

CDC automático bitemporal

Importante

AUTO CDC bitemporal está en fase Beta.

ScD Tipo 1 y Tipo 2 son unitemporales: realizan un seguimiento de los cambios en una sola dimensión de tiempo. Bitemporal amplía el historial de SCD de tipo 2 para hacer un seguimiento de los cambios en dos dimensiones temporales y distinguir entre dos perspectivas:

  • Tiempo de negocio: cuando se produjo realmente el evento.
  • Hora del sistema: cuando el sistema registró o ingerió el evento.

Al igual que el Tipo 2 de SCD, bitemporal preserva un historial completo de los registros. Agrega una segunda escala de tiempo para que pueda reconstruir tanto lo que mostraron los datos como lo que el sistema creía en cualquier punto del pasado.

Por ejemplo, un fondo de cobertura incorpora datos bursátiles de un sistema de origen. El precio de acciones de Acme Corp cambia el 1 de enero, pero el fondo no ingiere esa actualización hasta el 5 de enero. Bitemporales AUTO CDC permite al fondo responder a dos preguntas distintas: cuál era el precio real de acciones de Acme Corp el 1 de enero (tiempo de negocio) y qué precio creía el sistema cuando el fondo tomó decisiones comerciales el 3 de enero (tiempo del sistema). La capacidad de distinguir entre estas escalas de tiempo es útil para la auditoría, los informes normativos y la toma de decisiones financieras.

Para habilitar el procesamiento bitemporales, establezca STORED AS BITEMPORAL (SQL) o stored_as_scd_type="bitemporal" (Python), use SEQUENCE BY para la columna de hora del negocio y use SYSTEM SEQUENCE BY para la columna de hora del sistema. La tabla de destino agrega columnas __SYSTEM_START_AT y __SYSTEM_END_AT junto con las columnas __START_AT y __END_AT de tipo 2 de SCD. Para más información sobre la sintaxis, consulte AUTO CDC INTO (pipelines) o create_auto_cdc_flow.

Ejemplos de AUTO CDC bitemporales

En el ejemplo siguiente se crea una tabla de destino bitemporal a partir de un pequeño conjunto de eventos CDC sintéticos. La bt columna tiene la hora empresarial y la st columna tiene la hora del sistema.

Python

from pyspark import pipelines as dp

# Source: synthetic CDC events
dp.create_streaming_table(name="cdc_source")

@dp.append_flow(target="cdc_source", once=True)
def load_cdc_source():
  return spark.createDataFrame(
    [
      (1, "x10", "y10", 10, 100),
      (1, "x20", "y20", 20, 200)
    ],
    schema="id INT, x STRING, y STRING, bt INT, st INT",
  )

# Target: bitemporal table
dp.create_streaming_table(name="target_bitemporal")

dp.create_auto_cdc_flow(
  target = "target_bitemporal",
  source = "cdc_source",
  keys = ["id"],
  sequence_by = "bt",
  system_sequence_by = "st",
  stored_as_scd_type = "bitemporal"
)

SQL

-- Source: synthetic CDC events
CREATE OR REFRESH STREAMING TABLE cdc_source_sql;

CREATE FLOW cdc_source_sql AS INSERT INTO ONCE
  cdc_source_sql BY NAME
SELECT * FROM VALUES
  (1, 'x10', 'y10', 10, 100),
  (1, 'x20', 'y20', 20, 200)
  AS t(id, x, y, bt, st);

-- Target: bitemporal table
CREATE OR REFRESH STREAMING TABLE target_bitemporal_sql;

CREATE FLOW target_bitemporal_sql AS AUTO CDC INTO
  target_bitemporal_sql
FROM
  stream(cdc_source_sql)
KEYS
  (id)
SEQUENCE BY
  bt
SYSTEM SEQUENCE BY
  st
STORED AS
  BITEMPORAL;

En la siguiente secuencia de cambios se muestra cómo una tabla bitemporal registra una inserción, una actualización, una actualización fuera de orden y una eliminación para una sola empresa. La columna de secuenciación genera las __START_AT columnas y __END_AT (tiempo de negocio) y la columna de secuenciación del sistema genera las __SYSTEM_START_AT columnas y __SYSTEM_END_AT (hora del sistema):

Columna Description
__START_AT Fecha y hora de negocio en que esta fila pasó a ser válida.
__END_AT Hora de negocio en la que finaliza la validez de esta fila. null si es válido indefinidamente.
__SYSTEM_START_AT El momento del sistema en el que se sabe que los datos de esta fila y su intervalo de tiempo de negocio son válidos.
__SYSTEM_END_AT La hora del sistema en la que se sabe que los datos de esta fila y su intervalo de tiempo de negocio quedan invalidados. null si se sabe que será verdadero indefinidamente.

El sistema controla los eventos que llegan en cualquier orden en ambas escalas de tiempo. Cuando un evento llega con un tiempo de negocio o un tiempo del sistema anterior al de los eventos ya procesados, el sistema corrige el historial afectado en lugar de limitarse a añadirlo únicamente al final.

Cambio 1: Insertar

La compañía A se añade el 18/07/2025 a las 10:01:00 (hora de negocio), pero no se incorpora al sistema hasta las 10:05:00 (hora del sistema).

Entrada:

CompanyId Punto de datos Secuenciación Secuenciación del sistema Operation
A XFv1 7/18/2025 10:01:00 7/18/2025 10:05:00 INSERT

Salida:

CompanyId Punto de datos __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULO 7/18/2025 10:05:00 NULO

XFv1 es válido a partir de las 10:01:00 sin fin conocido. El sistema tuvo conocimiento de este hecho a la hora del sistema 10:05:00, sin hora de finalización conocida.

Cambio 2: Actualización

La empresa A se actualizó el 18/7/2025 a las 12:15:43 (hora de negocio), y el sistema procesa el evento a las 12:20:00 (hora del sistema). El sistema conserva lo que creía antes de que se conocía la actualización y el historial de negocios corregido después de ingerir la actualización.

Entrada:

CompanyId Punto de datos Secuenciación Secuenciación del sistema Operation
A XFv2 7/18/2025 12:15:43 7/18/2025 12:20:00 UPDATE

Salida:

CompanyId Punto de datos __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULO 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 NULO
A XFv2 7/18/2025 12:15:43 NULO 7/18/2025 12:20:00 NULO

XFv1 fue considerado válido de 10:01:00 sin fin conocido, y el sistema mantuvo esa creencia de 10:05:00 hasta las 12:20:00. Ahora se sabe que XFv1 solo es válido hasta las 12:15:43; un historial corregido entra en vigor a partir de la hora del sistema 12:20:00, sin una hora de finalización conocida. XFv2 es válido a partir de las 12:15:43 sin un final conocido y se aprendió en la hora del sistema 12:20:00.

Cambio 3: Actualización desordenda

Llega una actualización fuera de orden que indica que la Compañía A se actualizó realmente a las 12:05:00 del 18/7/2025 (hora de negocio), pero no se incorpora hasta las 12:25:00 (hora del sistema). Cuando una actualización llega más tarde según el tiempo del sistema, pero con un tiempo de negocio anterior, el sistema corrige el tiempo de negocio histórico y conserva tanto lo que el sistema creía antes de la actualización fuera de orden como el historial corregido.

Entrada:

CompanyId Punto de datos Secuenciación Secuenciación del sistema Operation
A XFv3 7/18/2025 12:05:00 7/18/2025 12:25:00 UPDATE

Salida:

CompanyId Punto de datos __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULO 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 NULO
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 NULO
A XFv2 7/18/2025 12:15:43 NULO 7/18/2025 12:20:00 NULO

Se consideró que XFv1 era válido desde las 10:01:00 hasta las 12:15:43, y esa consideración es ahora válida según el tiempo del sistema hasta las 12:25:00. La nueva actualización corrige la vigencia de negocio de XFv1 para que finalice a las 12:05:00, con un historial corregido efectivo desde el tiempo del sistema 12:25:00. Ahora se sabe que XFv3 es válido desde las 12:05:00 hasta las 12:15:43, una afirmación válida en el tiempo del sistema desde las 12:25:00 y sin final conocido.

Cambio 4: Eliminar

La empresa A se elimina a las 18/7/2025 12:30:00 y el sistema consume el evento a las 12:30:00. Dado que una operación de eliminación representa el final de la existencia empresarial de la entidad, el sistema no crea ninguna fila de reemplazo. XFv2 aparece en dos filas, conservando un rastro de auditoría completo de cuándo la empresa dejó de existir y de cuándo el sistema tuvo constancia de la eliminación.

Entrada:

CompanyId Punto de datos Secuenciación Secuenciación del sistema Operation
A XFv2 7/18/2025 12:30:00 7/18/2025 12:30:00 DELETE

Salida:

CompanyId Punto de datos __START_AT __END_AT __SYSTEM_START_AT __SYSTEM_END_AT
A XFv1 7/18/2025 10:01:00 NULO 7/18/2025 10:05:00 7/18/2025 12:20:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:15:43 7/18/2025 12:20:00 7/18/2025 12:25:00
A XFv1 7/18/2025 10:01:00 7/18/2025 12:05:00 7/18/2025 12:25:00 NULO
A XFv3 7/18/2025 12:05:00 7/18/2025 12:15:43 7/18/2025 12:25:00 NULO
A XFv2 7/18/2025 12:15:43 NULO 7/18/2025 12:20:00 7/18/2025 12:30:00
A XFv2 7/18/2025 12:15:43 7/18/2025 12:30:00 7/18/2025 12:30:00 NULO

XFv2 fue válido desde las 12:15:43, sin un final conocido, y el sistema mantuvo esa creencia desde las 12:20:00 hasta las 12:30:00. Una vez ingerida la eliminación, se sabe que XFv2 solo es válido hasta las 12:30:00 y corresponde a un historial corregido vigente desde la hora del sistema 12:30:00.

¿Qué objetos de datos se utilizan para el procesamiento de CDC en un pipeline?

Al declarar la tabla de destino en el metastore de Hive, se crean dos estructuras de datos:

  • Vista con el nombre asignado a la tabla de destino.
  • Una tabla interna de respaldo utilizada por el pipeline para administrar el procesamiento de CDC. Esta tabla se denomina anteponiendo __apply_changes_storage_ al nombre de la tabla de destino.

Por ejemplo, si declara una tabla de destino denominada dp_cdc_target, verá una vista denominada dp_cdc_target y una tabla denominada __apply_changes_storage_dp_cdc_target en el metastore. Consulte la vista para acceder a los datos procesados. No modifique directamente la tabla de respaldo.

Nota:

Estas estructuras de datos se aplican solo al procesamiento de AUTO CDC, no al procesamiento de AUTO CDC FROM SNAPSHOT. También se aplican solo a la tienda de metadatos de Hive, no al catálogo de Unity.