Nota:
El acceso a esta página requiere autorización. Puede intentar iniciar sesión o cambiar directorios.
El acceso a esta página requiere autorización. Puede intentar cambiar los directorios.
Importante
Esta característica está en versión preliminar pública. Los administradores del área de trabajo pueden controlar el acceso a esta característica desde la página Vistas previas . Consulte Administrar versiones preliminares de Azure Databricks.
Un objeto Stream representa un origen de datos de streaming externo, como Apache Kafka. Los flujos almacenan detalles de conexión, autenticación, esquemas y configuración de ingesta. Una vez creado un flujo, puede hacer referencia a él mediante definiciones de Feature View para crear funciones de streaming en tiempo real.
Los flujos tienen nombres de tres partes (catalog.schema.stream_name). El acceso a un flujo se rige por su tabla de ingestión asociada. Consulta Ingestión y relleno para más detalles.
Requisitos
- Para ejecutar comandos de cuaderno: sin servidor o un clúster de proceso clásico que ejecuta Databricks Runtime 17.0 ML o superior.
- Se debe instalar el paquete de Python
feature-engineering-clientversión 0.18.0 o superior.
Conectarse a fuentes de streaming
Antes de definir las funciones de streaming, conecte y pruebe la conexión de una canalización de Lakeflow Streaming con su broker de Kafka. Feature Store depende de un SDP serverless, lo que significa que necesitarás un mecanismo para conectar tu computación clásica (broker o endpoint) con la computación serverless de Databrick. Esto se hace a través de productos como privatelink o permitiendo que tu computación clásica sea accesible desde internet público.
Crear un flujo
Usa create_stream() para crear un nuevo Stream. Un Stream requiere cuatro componentes de configuración:
- Configuración de fuente: Especifica la plataforma de streaming y detalles específicos de la fuente, como la suscripción al tema de una fuente Kafka.
- Configuración de conexión: especifica cómo conectarse y autenticarse en la plataforma de streaming, incluidos los servidores de arranque y las credenciales.
- Configuración del esquema: define la estructura de las claves y los valores del mensaje.
- Configuración de ingesta: especifica dónde y cómo se ingieren los datos de flujo. Consulta Ingestión y relleno para más detalles.
Para la configuración específica source_config de la fuente y la conexión, junto con un ejemplo completo create_stream() , véase Apache Kafka. Las opciones de esquema y de ingestión se comparten entre fuentes.
Apache Kafka
Para transmitir desde Apache Kafka, usa KafkaStreamConfig como configuración de origen y una conexión de Unity Catalog para autenticación. Consulte Streaming en proceso sin servidor y Conectarse a Apache Kafka para obtener información sobre la conectividad con Kafka.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
DirectSchemas,
SchemaConfig,
IngestionConfig,
IngestionDestination,
StreamBackfillSource,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="events-topic"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "transaction_id": {"type": "string"},'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string", "format": "date-time"}'
' }'
'}'
)
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
),
)
Modos de suscripción Kafka
El modo de suscripción especifica cómo Stream selecciona los topics de Kafka de los que consume. Se admiten tres modos:
| Modo | Description | Ejemplo |
|---|---|---|
subscribe |
Lista separada por comas de nombres de temas | KafkaSubscriptionMode(subscribe="topic1,topic2") |
subscribe_pattern |
Coincidencia de nombres de tema con patrones de expresiones regulares en Java | KafkaSubscriptionMode(subscribe_pattern="events-.*") |
assign |
JSON que especifica asignaciones de particiones de tema | KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}') |
Autenticación de Kafka
Conexión de Unity Catalog (recomendado)
Use una conexión de Catálogo de Unity para autenticarse en el clúster de Kafka. Este es el enfoque recomendado para la autenticación administrada. Para crear una conexión, consulte Creación de una conexión. El creador del flujo debe tener USE CONNECTION en la conexión. Cualquier usuario que cree funcionalidades a partir de transmisiones como origen también debe tener USE CONNECTION en la conexión.
connection_config = StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
)
La conexión soporta tanto autenticación IAM (credencial de servicio) como SASL.
IAM (credencial de servicio)
Autentica con una credencial de servicio de Unity Catalog, por ejemplo, para conectarte a Amazon MSK con IAM. Para crear una credencial de servicio, consulte Creación de credenciales de servicio. Establece el nombre de la credencial del servicio con la credential opción:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>'
)
Además de USE CONNECTION en la conexión, las identidades que utilizan la credencial de servicio necesitan ACCESS en ella. Conceda ACCESS en la credencial de servicio indicada al creador de la transmisión y a cualquier identidad que materialice características de la transmisión. Consulte Concesión de permisos para usar una credencial de servicio para acceder a un servicio en la nube externo.
SASL
La autenticación SASL utiliza un nombre de usuario y una contraseña. Establezca en sasl_mechanism una de las siguientes opciones:
PLAINSCRAM-SHA-256SCRAM-SHA-512
Proporciona las credenciales con las opciones user y password. La conexión almacena estas credenciales de forma segura.
El siguiente ejemplo utiliza SASL/SCRAM. Para SASL/PLAIN, pon sasl_mechanism en PLAIN.
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
sasl_mechanism 'SCRAM-SHA-512',
user '<username>',
password '<password>'
)
MTLS directo
Para la autenticación mTLS directa, proporcione archivos keystore y truststore almacenados en un volumen de Unity Catalog, con contraseñas referenciadas mediante secret scopes de Databricks. Para más información sobre la autenticación SSL con Kafka, consulte Uso de SSL para conectarse Azure Databricks a Kafka.
from databricks.feature_engineering.entities import (
DirectMtlsConfig,
MtlsConfig,
SecretScopeReference,
)
connection_config = DirectMtlsConfig(
bootstrap_servers="broker1:9092,broker2:9092",
mtls_config=MtlsConfig(
keystore_location="/Volumes/my_catalog/my_schema/my_volume/keystore.jks",
keystore_password_ref=SecretScopeReference(
scope="my_scope", key="keystore_password"
),
key_password_ref=SecretScopeReference(
scope="my_scope", key="key_password"
),
truststore_location="/Volumes/my_catalog/my_schema/my_volume/truststore.jks",
truststore_password_ref=SecretScopeReference(
scope="my_scope", key="truststore_password"
),
),
)
Configuración del esquema
Define la estructura de las claves y valores de los mensajes para que las definiciones de ingestión y características puedan leer campos individuales. Para los orígenes de Kafka, payload_schema corresponde al valor del mensaje de Kafka (en el value modelo clave-valor de Kafka) y key_schema corresponde a la clave de mensaje de Kafka. Se debe proporcionar al menos uno de payload_schema o key_schema .
Cada uno SchemaConfig acepta uno de tres formatos, coincidiendo con cómo la fuente serializa sus mensajes: json_schema, avro_schema, o proto_schema. Si no se proporciona ningún esquema para una clave o carga, se trata como una cadena simple.
Los ejemplos de código de esta sección utilizan esquemas declarados en línea con DirectSchemas, donde el esquema se proporciona como una cadena. Para gestionar esquemas utilizando un registro externo de esquemas, consulte Registro de esquemas para más detalles.
Esquema JSON
Proporciona una cadena de esquema JSON a json_schema.
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"}'
' }'
'}'
)
),
key_schema=SchemaConfig(
json_schema='{"type": "string"}'
),
)
Esquema de Avro
Proporciona una cadena de esquema Avro a avro_schema. Se admiten tipos lógicos Avro, incluyendo timestamp-millis, date, y decimal.
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
avro_schema=(
'{'
' "type": "record",'
' "name": "Event",'
' "fields": ['
' {"name": "user_id", "type": "string"},'
' {"name": "amount", "type": "double"},'
' {"name": "event_time",'
' "type": {"type": "long", "logicalType": "timestamp-millis"}}'
' ]'
'}'
)
),
)
Esquema Protobuf
Proporcione un ProtoSchemaSpec a proto_schema con el texto de origen de Búferes de protocolo.proto y el nombre del mensaje de la carga. Importa ProtoSchemaSpec desde databricks.feature_engineering.entities.
message_name debe ser el nombre del mensaje completamente cualificado, incluyendo el package declarado en el .proto texto (por ejemplo, com.example.Event, no Event). Se admiten tanto la sintaxis proto2 como la proto3.
google.protobuf.Timestamp y los tipos de envoltorios escalares (StringValue, Int32Value, y así sucesivamente) son soportados, y sus importaciones se resuelven automáticamente. Otros tipos conocidos, como Duration, Struct, y Any, son rechazados; codifica esos valores como un escalar o mensaje soportado en su lugar. Los tipos escalares fixed32 y fixed64, así como map con claves que no sean de cadena, tampoco son compatibles.
from databricks.feature_engineering.entities import ProtoSchemaSpec
schema_config = DirectSchemas(
payload_schema=SchemaConfig(
proto_schema=ProtoSchemaSpec(
schema_text=(
'syntax = "proto3";\n'
'package com.example;\n'
'import "google/protobuf/timestamp.proto";\n'
'message Event {\n'
' string user_id = 1;\n'
' double amount = 2;\n'
' google.protobuf.Timestamp event_time = 3;\n'
'}'
),
message_name="com.example.Event",
)
),
)
Decodificación de datos usando esquemas
Databricks decodifica cada mensaje con las funciones de from_jsonSpark , from_avro, y from_protobuf . Se aplican los siguientes comportamientos tanto si declaras el esquema en línea como si lo resuelves desde un registro de esquemas:
- Registros malformados. La decodificación utiliza el modo
PERMISSIVE, por lo que un registro que no coincide con su esquema da como resultado un valor nulo en lugar de provocar un error en el flujo. - Sindicatos Avro. Una unión de múltiples tipos de registro se decodifica como una estructura con un campo por cada tipo de registro, y cada campo recibe el nombre de su registro de Avro correspondiente.
- Tipos Protobuf. Los enteros sin signo se decodifican como un tipo con signo de mayor tamaño (por ejemplo,
uint32aBIGINTyuint64aDECIMAL(20,0)), los campos de enumeración se decodifican como su nombre en forma de cadena, y los tipos envoltorio escalares (por ejemplo,StringValueyInt32Value) se decodifican como una columna que admite valores NULL del tipo envuelto.
Registro de esquemas
Los registros de esquemas almacenan y versionan los esquemas que usan los productores y consumidores de streaming, haciendo cumplir las reglas de compatibilidad a medida que esos esquemas evolucionan. Cuando se configura un registro de esquema externo, Feature Store lee el esquema del registro y lo utiliza para decodificar el mensaje en streaming. No declaras el esquema en línea en el Stream cuando usas un registro de esquema.
El soporte para registros de esquemas tiene las siguientes limitaciones:
- Compatible únicamente con flujos de Kafka.
- Solo es compatible Confluent Schema Registry
- Solo se soportan los formatos Avro y Protobuf . Para leer mensajes JSON, declara el esquema en línea en su lugar. Ver esquema JSON.
- Cada Flujo está conectado exactamente a un sujeto Confluente para el valor del mensaje y uno para la clave del mensaje (si se proporciona). Los temas de flujo que contienen múltiples registros de esquema no son una configuración soportada. Si tu Stream se conecta a temas que contienen múltiples esquemas, los registros que no coinciden con el esquema del tema especificado se decodifican como nulos.
Conectarse a un registro de esquemas
Proporcione los detalles de la conexión del registro como opciones en la conexión del Catálogo Kafka Unity y almacene el secreto de la API del registro en un ámbito secreto de Databricks. La identidad de ejecución de transmisión debe tener permiso READ en el ámbito de secretos, porque la canalización de ingesta lee el secreto en tiempo de ejecución. Para saber cómo crear y configurar una conexión, véase Crear una conexión.
Añade las schema_registry_urlopciones , schema_registry_api_key, y schema_registry_api_secret a la conexión utilizada para la autenticación. El siguiente ejemplo crea una conexión Kafka que se autentica con el broker con una credencial de servicio de Unity Catalog y con el registro con una clave API:
CREATE CONNECTION IF NOT EXISTS `my-kafka-connection`
TYPE KAFKA
OPTIONS (
bootstrap_servers '<bootstrap_servers>',
credential '<service_credential>',
schema_registry_url 'https://<registry-host>',
schema_registry_api_key '<registry_api_key>',
schema_registry_api_secret secret('<scope>', '<key>')
)
Establezca tanto la opción schema_registry_api_secret en la conexión de Kafka como la referencia al ámbito de secreto en el flujo para que apunten al mismo secreto.
Crea un flujo que utilice un registro de esquema
Pase a SchemaRegistryConfig como el schema_config. Referencia el secreto de la API del registro con api_secret_ref, e identifica el asunto y formato con payload_schema_locator para el valor del mensaje, o key_schema_locator para la clave del mensaje. Debe proporcionarse al menos un localizador.
Observa las diferencias aquí en comparación con los ejemplos directos de esquemas en la sección de configuración de esquemas . Al usar un registro de esquema, no se proporciona el esquema en línea en la transmisión a schema_config. En su lugar, especificas una SchemaRegistryConfig que identifica el esquema en el registro.
from databricks.feature_engineering import FeatureEngineeringClient
from databricks.feature_engineering.entities import (
KafkaStreamConfig,
KafkaSubscriptionMode,
StreamConnectionConfig,
SchemaRegistryConfig,
SchemaLocator,
SchemaLocatorConfluentSchema,
SchemaLocatorFormat,
SecretScopeReference,
IngestionConfig,
IngestionDestination,
)
client = FeatureEngineeringClient()
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
source_config=KafkaStreamConfig(
subscription_mode=KafkaSubscriptionMode(subscribe="transactions"),
),
connection_config=StreamConnectionConfig(
uc_connection_name="my-kafka-connection"
),
schema_config=SchemaRegistryConfig(
api_secret_ref=SecretScopeReference(
scope="my_scope", key="sr_api_secret"
),
payload_schema_locator=SchemaLocator(
confluent_schema=SchemaLocatorConfluentSchema(
subject="transactions-value"
),
format=SchemaLocatorFormat.FORMAT_AVRO,
),
),
ingestion_config=IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.transactions_ingestion"
),
),
)
Un sujeto Confluente es el ámbito denominado bajo el cual se registra el historial de versiones de un esquema y se hace cumplir la compatibilidad. Establezca subject el nombre del ámbito correspondiente, que de forma común se determina a partir de la estrategia de nombres de sujeto:
-
TopicNameStrategy (por defecto, deriva el tema del nombre del tema):
<topic>-valuepara el valor y<topic>-keypara la clave. Por ejemplo, el esquema de valores del tematransactionsutiliza el sujetotransactions-value. -
RecordNameStrategy (deriva el sujeto del nombre del registro del esquema, independientemente del tema): el nombre de registro completamente cualificado, como
com.example.Payment. Este es el espacio de nombres y el nombre del registro para Avro, o el paquete y el nombre del mensaje para Protobuf. -
TopicRecordNameStrategy (combina los nombres de tema y registro):
<topic>-<fully-qualified-record-name>, comotransactions-com.example.Payment.
format es obligatorio. Configúralo para SchemaLocatorFormat.FORMAT_AVRO o SchemaLocatorFormat.FORMAT_PROTOBUF para que coincida con cómo se serializa el tema.
Evolución de esquemas
La canalización de ingesta resuelve el esquema actual del sujeto cuando se inicia. Cuando registra una nueva versión de esquema compatible con versiones anteriores en el sujeto en el registro de esquemas, la canalización en ejecución sigue usando la versión con la que empezó.
Para los Flujos respaldados por el registro de esquemas, la tubería de ingestión se reinicia automáticamente cada dos o tres horas. En cada reinicio, recupera la versión más reciente del esquema del subject, y los campos nuevos o modificados aparecen en la tabla de ingestión.
Los flujos que usan esquemas directos en lugar de un registro de esquemas evolucionan su esquema con update_stream. Ver Actualizar un flujo.
Para ver cómo la pipeline gestiona los registros que no coinciden con el esquema que está usando actualmente, véase Decodificación de datos usando esquemas.
Filtrar registros por tipo
Un flujo decodifica cada registro usando un único esquema de clave y valor (si se proporciona), ya sea que los especifiques directamente o utilices un registro de esquemas. Como un tópico puede contener más de un tipo de registro y los Streams pueden suscribirse a varios tópicos, use record_type_filter para seleccionar qué registros del tópico pertenecen a este Stream.
Proporciona una expresión SQL que haga referencia a los campos decodificados mediante notación de punto, por ejemplo value.event_type = 'transaction'. Los registros que no coinciden con el filtro se ignoran. No se escriben en la tabla de ingestión ni se usan en materialización. Para crear un flujo para otros tipos de registro, crea un flujo separado con otro record_type_filter.
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
record_type_filter="value.event_type = 'transaction'",
)
Incluso sin record_type_filter, la decodificación nunca falla al Stream. Un registro que no coincide con el esquema configurado se decodifica de forma permisiva. Los registros se decodifican de una de las siguientes maneras:
- En una fila con
NULLlos valores de los campos que el esquema espera pero que el registro omite (JSON, Avro y Protobuf). - En una fila que contiene valores que pertenecen a otro tipo de registro (solo Avro y Protobuf).
Para identificar si una fila en la tabla de ingestión pertenece al tipo de registro esperado, utilice una de las siguientes comprobaciones:
- Comprueba que un campo sea igual a un valor esperado,
value.event_type = 'transaction'por ejemplo (preferido para Avro y Protobuf). - Comprueba que un campo no sea
NULL, por ejemplovalue.activity_id IS NOT NULL.
Se recomienda usar record_type_filter con flujos separados cuando los esquemas difieren sustancialmente entre tipos de registro sobre el tema o cuando se quiere gobernar el acceso a cada tipo de registro de forma independiente. Para mantener los costes, Databricks recomienda mantener un pequeño número de Flujos, ya que cada Flujo tiene una canalización de ingestión y una tabla de ingestión separadas. Cada Stream también utiliza recursos de cómputo independientes en el momento de la materialización. Puedes usar filtros específicos de cada función para la materialización.
record_type_filter es diferente de la filter_conditioncaracterística .
record_type_filter se configura en el Stream y controla qué registros se ingieren y están disponibles para todas las funciones que utilizan el Stream como origen, mientras que filter_condition se configura en una función individual y filtra las filas antes de la agregación. Consulta Condiciones de filtro en fuentes de streaming para más detalles sobre filter_condition.
Ingestión y relleno
El ingestion_config parámetro configura cómo se capturan y almacenan los datos de flujo para el entrenamiento y el servicio.
El acceso a un Stream se rige por la tabla de ingestión:
-
SELECTen la tabla de ingesta otorga acceso de lectura al Stream. -
MANAGEen la tabla de ingesta concede acceso para eliminar.
Para obtener más información sobre los privilegios de tabla, consulte Table y la referencia de privilegios de Unity Catalog.
Canalización de ingesta
Cuando se crea un flujo, Databricks inicia una tubería de ingestión gestionada que lee continuamente los mensajes del flujo fuente y los escribe en una tabla Delta (la tabla de ingestión). La canalización se inicia en la posición más reciente del origen y se ejecuta de forma continua, capturando solo los mensajes nuevos que llegan después de que se crea el flujo. Esta tabla de ingesta se usa para el entrenamiento con características de streaming. Cuando se elimina un flujo, también se eliminan la canalización de ingesta y la tabla de ingesta.
Destino de la ingesta
ingestion_destination especifica el nombre de la tabla Delta de tres partes donde se escriben los datos del flujo.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
)
Esquema de tabla de ingesta
La tabla de ingestión contiene los datos del mensaje junto con las columnas de metadatos. Las columnas comunes se encuentran en cada fuente; las columnas kafka_* se encuentran solo en un flujo de Kafka.
| Columna | Tipo | Source | Description |
|---|---|---|---|
key |
Varía (según key_schema) |
Common | La clave del mensaje, estructurada según el esquema que proporcionaste. |
value |
Varía (según payload_schema) |
Common | El valor del mensaje (payload), estructurado según el esquema que proporcionaste. |
stream_record_timestamp |
TIMESTAMP |
Common | Marca de tiempo del registro. Para los datos rellenados hacia adelante, esta es la marca de tiempo de ingesta del origen. En el caso de los datos de relleno, los proporciona el cliente. |
record_source |
STRING |
Common | Ya sea "stream" (relleno directo desde la transmisión en directo) o "backfill" (desde la fuente de relleno). |
kafka_topic |
STRING |
Kafka | El tema kafka del que se consumió el registro. |
kafka_partition |
INT |
Kafka | Partición de Kafka desde la que se consumió el registro. |
kafka_offset |
LONG |
Kafka | El desfase de Kafka del registro dentro de su partición. |
Origen de reposición
Dado que la canalización de llenado comienza a partir de la última posición del origen, no recoge los mensajes que existían antes de que se creara la transmisión. Para proporcionar cobertura histórica de datos para el entrenamiento, configure un origen de reposición opcional.
Cuando se configura un origen de reposición, Databricks ejecuta un trabajo único de MERGE INTO que copia las filas de reposición en la tabla de ingesta con record_source="backfill". MERGE solo se ejecuta después de que el comprobador de reposición confirme que la fuente de backfill y el flujo de relleno hacia delante tienen marcas de tiempo superpuestas (consulte Superposición entre los datos de backfill y el flujo en vivo). Si no se cumple la condición de superposición en un plazo de 2 días, MERGE se ejecuta de todos modos para evitar el bloqueo indefinidamente.
La tabla de reposición debe incluir una stream_record_timestamp columna de tipo TIMESTAMP en la zona horaria UTC. Las demás columnas de metadatos se transfieren si están presentes en la fuente de backfill o, en caso contrario, se establecen en NULL. Para Kafka, estos son kafka_topic, kafka_partition, y kafka_offset.
from databricks.feature_engineering.entities import StreamBackfillSource
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
backfill_source=StreamBackfillSource(
delta_table_name="my_catalog.my_schema.historical_events"
),
)
Superposición entre los datos de reposición y streaming en vivo
Antes de ejecutar una combinación entre reposición y la tabla de ingesta, una comprobación de superposición compara las marcas de tiempo en las dos tablas:
-
Número máximo de reposición: el máximo
stream_record_timestampen el origen de reposición. -
Mín. de ingesta: el mínimo
stream_record_timestampde filas (record_source="stream") en la tabla de ingesta.
La operación MERGE se lleva a cabo cuando la marca de tiempo más reciente de la carga de reposición supera en al menos 1 hora la marca de tiempo más antigua de la tabla de ingesta. Esta superposición garantiza que no haya huecos en la tabla de ingesta. Si no se cumple la condición de superposición en un plazo de 2 días, MERGE se ejecuta de todos modos para evitar el bloqueo indefinidamente.
Como la tubería de ingestión comienza desde la posición más reciente en la fuente, solo captura los mensajes que llegan después de que se crea el flujo. El origen de reposición debe contener datos que se extienden al intervalo de tiempo de ingesta, no solo hasta el tiempo de creación de la transmisión.
Por ejemplo, si crea una transmisión a las 3:00 p. m., la pipeline de reposición comienza a leer mensajes a partir de las 3:00 p. m. El origen de reposición debe incluir datos con marcas temporales hasta al menos las 4:00 p. m. (1 hora después del inicio de la reposición) para superar la comprobación de solapamiento. Esto significa que debe actualizar la tabla de reposición después de las 4:00 p. m. para asegurarse de que la tabla de ingesta no tiene huecos.
Desduplicación
Utilice deduplication_columns para especificar rutas de columnas a fin de identificar filas duplicadas durante la ingesta de datos entre los flujos de relleno retroactivo y de relleno progresivo. Use la notación de puntos para los campos anidados (por ejemplo, "value.user_id").
Elija columnas de desduplicación en función de los datos:
- Si cada registro de la secuencia contiene un identificador único (por ejemplo,
value.transaction_id), use esa columna para la desduplicación. - Si el origen de reposición incluye
kafka_partitionykafka_offsetcolumnas, úselos para identificar de forma única cada registro. - Si no se especifica ninguna columna de desduplicación, la clave de desduplicación predeterminada es la combinación completa de
key,valueystream_record_timestamp. Esto no se recomienda, ya que esta coincidencia estricta de criterios puede provocar fácilmente duplicados.
ingestion_config = IngestionConfig(
ingestion_destination=IngestionDestination(
delta_table_name="my_catalog.my_schema.events_ingestion"
),
deduplication_columns=["value.transaction_id"],
)
Atribución de costos
Establezca tags y budget_policy_id en IngestionConfig para atribuir el coste de la ingesta gestionada del flujo. Azure Databricks los aplica a la canalización de ingesta de Lakeflow y a sus trabajos de relleno hacia delante y relleno retrospectivo cuando se crea el flujo.
Por ejemplo, los límites de etiquetas y cómo consultar el gasto atribuido, véase Costes de atributos con etiquetas y políticas de uso serverless.
Excluir columnas de un flujo
Usa excluded_columns para eliminar columnas específicas de un Stream que no quieres ingerir. Una columna excluida no se escribe en la tabla de ingesta y no se puede hacer referencia a ella mediante un atributo ni utilizarse en el entrenamiento.
Especifica cada columna usando notación de puntos en la clave o valor del mensaje, como value.user.email o key.account_id. Estas columnas se colocan desde el decodificado key y value a través de la ingestión, relleno y materialización. Si una ruta apunta a una estructura, también se eliminan todos sus campos anidados (por ejemplo, al eliminar value.address también se eliminan value.address.city y value.address.zip).
stream = client.create_stream(
name="my_catalog.my_schema.my_stream",
# ...source, connection, schema, and ingestion config...
excluded_columns=["value.user.email", "value.user.ssn"],
)
Al usar esquemas directos, la columna excluida debe existir ya en el esquema clave o de valores, o create_stream falla. Al usar un registro de esquemas, puedes excluir una columna antes de que exista. Una columna excluida tampoco puede ser una columna de deduplicación, porque las columnas de deduplicación son necesarias para identificar filas duplicadas. Cualquier característica que haga referencia a una columna excluida (por ejemplo, como entidad, serie temporal o entrada) no se crea.
Puedes cambiar las columnas excluidas de un Stream después de su creación con update_stream, tanto en Streams con esquema directo como en Streams respaldados por un registro de esquemas. Consulta Actualizar un stream para más detalles.
Gestionar flujos
Obtener una transmisión
stream = client.get_stream(name="my_catalog.my_schema.my_stream")
Listar flujos
streams = client.list_streams(
catalog_name="my_catalog",
schema_name="my_schema",
max_results=50,
include_schemas=False,
)
Establezca include_schemas=True para incluir todos los detalles del esquema. Los esquemas pueden ser grandes y esto puede dar lugar a una operación de larga duración. Para recuperar esquemas individualmente, use get_stream.
Actualizar un flujo
Usa update_stream para cambiar un flujo después de crearlo. Pase schema_config para hacer evolucionar un esquema directo, excluded_columns para cambiar qué columnas se descartan, o ambos. No se permite actualizar otros campos. Crea un nuevo flujo en su lugar.
Actualizar un flujo reinicia su pipeline de ingesta para que el cambio tenga efecto. La ingestión suele reanudarse en pocos minutos.
Desarrollar un esquema directo
Para un Stream que utiliza esquemas directos, pasa un DirectSchemas a schema_config. Establece payload_schema, key_schema, o ambos. Un bando que no marcas queda sin cambios. Los flujos respaldados por el registro de esquemas rechazan una schema_config actualización y deben evolucionarse a través del registro.
from databricks.feature_engineering.entities import DirectSchemas, SchemaConfig
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
schema_config=DirectSchemas(
payload_schema=SchemaConfig(
json_schema=(
'{'
' "type": "object",'
' "properties": {'
' "user_id": {"type": "string"},'
' "amount": {"type": "number"},'
' "event_time": {"type": "string"},'
' "channel": {"type": "string"}'
' }'
'}'
)
),
),
)
Las actualizaciones del esquema deben ser retrocompatibles para que la canalización de ingesta en ejecución pueda seguir descodificando los registros existentes y escribiendo en la tabla de ingesta. Cualquier otro cambio se rechaza.
Lo que está permitido depende del formato:
-
JSON y Protobuf: añadir campos opcionales, eliminar campos y ampliar el tipo de campo (por ejemplo,
intabigint). Protobuf también permite reordenar campos. -
Avro: solo permite ensanchar
intylongeliminar un campo final cuyos bytes no leen el campo posterior. Para evolucionar un esquema Avro de forma más libre, utiliza un flujo respaldado por un registro de esquemas.
Sumar campos hace crecer la decodificación key y value structs de la tabla de ingestión. Las filas escritas antes de la actualización mantienen su forma original, y los campos añadidos se leen igual que NULL las filas anteriores. Las eliminaciones y cambios de tipo solo entran en vigor para los registros ingeridos tras la actualización.
Columnas excluidas por cambio
Pasa el nuevo conjunto completo de rutas de columna a excluded_columns, que reemplaza el conjunto existente. Pasa una lista vacía ([]) para eliminar todas las exclusiones. Para más detalles sobre este comportamiento, véase Excluir columnas de un Flujo.
stream = client.update_stream(
name="my_catalog.my_schema.my_stream",
excluded_columns=["value.user.email", "value.user.ssn"],
)
Cambiar las columnas excluidas es solo hacia adelante. Las columnas recién excluidas dejan de escribirse (aparecen en NULL) y las nuevas incluidas empiezan a poblarse a partir de ahora, mientras que las filas previamente escritas quedan as-is. Para evitar que se ingiera una nueva columna:
-
Registro de esquemas: añade primero la columna
excluded_columnsy espera a que se reinicie la tubería de ingestión, luego registra la nueva versión del esquema en el registro. -
Esquema directo: añadir la columna a
schema_configy aexcluded_columnsen la mismaupdate_streamllamada.
Eliminar un flujo
Al eliminar un flujo, también se eliminan su canalización de ingesta y su tabla de ingesta.
Warning
Los modelos o características que hacen referencia a la secuencia eliminada ya no tendrán acceso a los datos de la secuencia subyacente. Cree una copia de la tabla de ingesta antes de eliminarla si necesita estos datos, pero ya no la transmisión.
client.delete_stream(name="my_catalog.my_schema.my_stream")
Cuaderno de ejemplo
Para ver un ejemplo de principio a fin que crea una transmisión, define las características de transmisión y despliega en un punto de conexión de servicio, consulte el siguiente cuaderno: