Configurare un flusso

Importante

Questa funzionalità è in Anteprima Pubblica. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.

Un oggetto Stream rappresenta un'origine dati di streaming esterna, ad esempio Apache Kafka. Streams archivia i dettagli della connessione, l'autenticazione, gli schemi e la configurazione di acquisizione. Dopo aver creato un flusso, è possibile farvi riferimento usando le definizioni di Visualizzazione funzionalità per creare funzionalità di streaming in tempo reale.

I flussi hanno nomi in tre parti (catalog.schema.stream_name). L'accesso a uno Stream è regolato dalla relativa tabella di acquisizione associata. Per informazioni dettagliate, vedere Inserimento e riempimento .

Requisiti

  • Per eseguire i comandi del notebook: serverless o un cluster di calcolo classico che esegue Databricks Runtime 17.0 ML o versione superiore.
  • È necessario installare il feature-engineering-client pacchetto Python versione 0.16.0 o successiva.

Creare un flusso

Usare create_stream() per creare un nuovo flusso. Un flusso richiede quattro componenti di configurazione:

  • Configurazione origine: specifica la piattaforma di streaming (ad esempio, Kafka) e i dettagli specifici dell'origine (ad esempio, la sottoscrizione dell'argomento per Kafka).
  • Configurazione connessione: specifica come connettersi ed eseguire l'autenticazione alla piattaforma di streaming, inclusi server e credenziali bootstrap.
  • Configurazione dello schema: definisce la struttura delle chiavi e dei valori dei messaggi.
  • Configurazione di acquisizione: specifica dove e come vengono acquisiti i dati del flusso. Per informazioni dettagliate, vedere Inserimento e riempimento .
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"
        ),
    ),
)

Connessione alle sorgenti di flusso

Prima di definire le funzionalità di streaming, stabilisci e testa la connessione di una pipeline Lakeflow Streaming al broker Kafka. Vedi Streaming sul calcolo serverless e Connettersi ad Apache Kafka.

Per il servizio di streaming gestito di AWS (Amazon MSK), consulta Connettività privata serverless verso Amazon MSK. Per informazioni dettagliate sulle opzioni di autenticazione Kafka, vedere Autenticazione.

Authentication

Usare una connessione del catalogo Unity per eseguire l'autenticazione al cluster Kafka. Questo è l'approccio consigliato per l'autenticazione gestita. Per creare una connessione, vedere Creare una connessione.

connection_config = StreamConnectionConfig(
    uc_connection_name="my-kafka-connection"
)

MTLS diretto

Per l'autenticazione mTLS diretta, specificare i file keystore e truststore archiviati in un volume di Unity Catalog, con le password referenziate tramite gli scope dei segreti di Databricks. Per altre informazioni sull'autenticazione SSL con Kafka, vedere Usare SSL per connettersi 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"
        ),
    ),
)

SASL

L'autenticazione SASL (SASL/SCRAM e SASL/PLAIN) non è supportata durante l'anteprima.

Modalità di sottoscrizione

La modalità di sottoscrizione specifica il modo in cui Stream seleziona gli argomenti Kafka da utilizzare. Sono supportate tre modalità:

Modalità Description Example
subscribe Elenco delimitato da virgole di nomi di argomenti KafkaSubscriptionMode(subscribe="topic1,topic2")
subscribe_pattern Java nomi di argomenti corrispondenti ai criteri regex KafkaSubscriptionMode(subscribe_pattern="events-.*")
assign JSON che specifica le assegnazioni di topic-partition KafkaSubscriptionMode(assign='{"my-topic": [0, 1, 2]}')

Configurazione dello schema

Definire la struttura delle chiavi e dei valori dei messaggi usando il formato dello schema JSON . Per le origini Kafka, payload_schema corrisponde al valore del messaggio Kafka (modello value chiave-valore di Kafka) e key_schema corrisponde alla chiave del messaggio Kafka. È necessario specificare almeno uno di payload_schema o key_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"}'
    ),
)

Se non viene fornito alcuno schema per una chiave o un payload, viene considerato come una stringa semplice.

Inserimento e riempimento

Il ingestion_config parametro consente di configurare la modalità di acquisizione e archiviazione dei dati di flusso per il training e la gestione.

L'accesso a un flusso è regolato dalla tabella di inserimento:

  • SELECT nella tabella di acquisizione concede l'accesso in lettura allo stream.
  • MANAGE nella tabella di inserimento concede l'accesso all'eliminazione.

Per ulteriori informazioni sui privilegi della tabella, vedere Tabella e riferimento ai privilegi di Unity Catalog.

Pipeline di inserimento

Quando viene creato un flusso, Databricks avvia una pipeline di inserimento gestita che legge continuamente i messaggi dall'argomento Kafka e li scrive in una tabella Delta (tabella di inserimento). La pipeline parte dall'offset Kafka più recente e viene eseguita continuamente, acquisendo solo i nuovi messaggi che arrivano dopo che il flusso è stato creato. Questa tabella di acquisizione viene utilizzata per l'addestramento con funzionalità in streaming. Quando un flusso viene eliminato, vengono eliminate anche la pipeline di inserimento e la tabella di inserimento.

Destinazione di acquisizione

ingestion_destination Specifica il nome della tabella Delta in tre parti in cui vengono scritti i dati del flusso.

ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
)

Schema della tabella di acquisizione

La tabella di inserimento contiene i dati del messaggio insieme alle colonne di metadati:

Column Digita Description
key Variabile (da key_schema) Chiave del messaggio Kafka, strutturata in base allo schema specificato.
value Variabile (da payload_schema) Valore del messaggio Kafka (payload), strutturato in base allo schema specificato.
stream_record_timestamp TIMESTAMP Data e ora del record. Per i dati di riempimento in avanti, si tratta del timestamp di inserimento del broker Kafka. Per i dati di backfill, questo è fornito dal cliente.
kafka_topic STRING L'argomento Kafka da cui è stato utilizzato il record.
kafka_partition INT La partizione Kafka da cui è stato utilizzato il record.
kafka_offset LONG L'offset Kafka del record all'interno della sua partizione.
record_source STRING O "stream" (riempimento in avanti dal flusso Kafka in tempo reale) o "backfill" (dall'origine di backfill).

Origine del riempimento retroattivo

Poiché la pipeline forward-fill inizia dall'offset Kafka più recente, non acquisisce i messaggi esistenti prima della creazione del flusso. Per fornire una copertura storica dei dati per l'addestramento, configura una fonte facoltativa di recupero dati pregressi.

Quando viene configurata un'origine di backfill, Databricks esegue un processo una tantum MERGE INTO che copia le righe di backfill nella tabella di ingestione con record_source="backfill". L'operazione MERGE viene eseguita solo dopo che la verifica della sovrapposizione conferma che l'origine del backfill e il flusso di forward-fill hanno timestamp sovrapposti (consulta Sovrapposizione tra il backfill e i dati del flusso live). Se la condizione di sovrapposizione non viene soddisfatta entro 2 giorni, l'operazione MERGE viene eseguita comunque per evitare il blocco illimitato.

La tabella backfill deve includere una stream_record_timestamp colonna di tipo TIMESTAMP nel fuso orario UTC. Le altre colonne di metadati di Kafka (kafka_topic, kafka_partition, kafka_offset) vengono propagate, se presenti nell'origine di backfill, oppure impostate su NULL in caso contrario.

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"
    ),
)

Sovrapposizione tra il riempimento retroattivo e i dati del flusso in tempo reale

Prima di eseguire un'operazione MERGE tra il backfill e la tabella di acquisizione, un controllo di sovrapposizione confronta i timestamp presenti nelle due tabelle:

  • Backfill max: valore massimo stream_record_timestamp nell'origine del riempimento.
  • Inserimento min: numero minimo stream_record_timestamp di righe (record_source="stream") nella tabella di inserimento.

Il MERGE procede quando il timestamp più recente del backfill supera di almeno 1 ora il timestamp più vecchio della tabella di ingestione. Questa sovrapposizione garantisce che non vi siano lacune nella tabella di inserimento. Se la condizione di sovrapposizione non viene soddisfatta entro 2 giorni, l'operazione MERGE viene eseguita comunque per evitare il blocco illimitato.

Poiché la pipeline di inserimento inizia dall'offset Kafka più recente, acquisisce solo i messaggi in arrivo dopo la creazione del flusso. L’origine del backfill deve contenere dati che si estendono nell’intervallo temporale di acquisizione, e non solo fino al momento di creazione del flusso.

Ad esempio, se si crea un flusso alle 15:00, la pipeline di forward-fill inizia a leggere i messaggi a partire dalle 15:00. L'origine dei dati di backfill deve includere dati con timestamp che arrivino almeno alle 16:00 (1 ora dopo l'inizio del forward-fill) per superare la verifica di sovrapposizione. Ciò significa che è necessario aggiornare la tabella backfill dopo le 14:00 per assicurarsi che la tabella di inserimento non contenga lacune.

Deduplicazione

Usare deduplication_columns per specificare i percorsi delle colonne per identificare le righe duplicate in fase di acquisizione tra i dati di flusso di backfill e forward-fill. Usare la notazione punto per i campi annidati , ad esempio "value.user_id".

Scegliere le colonne di deduplicazione in base ai dati:

  • Se ogni record nel flusso contiene un identificatore univoco ( ad esempio , value.transaction_id), usare tale colonna per la deduplicazione.
  • Se l'origine di backfill include le colonne kafka_partition e kafka_offset, usale per identificare in modo univoco ogni record.
  • Se non vengono specificate colonne di deduplicazione, la chiave di deduplicazione predefinita è la combinazione completa di key, valuee stream_record_timestamp. Questa operazione non è consigliata perché questa rigorosa corrispondenza dei criteri può causare facilmente duplicati.
ingestion_config = IngestionConfig(
    ingestion_destination=IngestionDestination(
        delta_table_name="my_catalog.my_schema.events_ingestion"
    ),
    deduplication_columns=["value.transaction_id"],
)

Gestire i flussi

Ottenere un flusso

stream = client.get_stream(name="my_catalog.my_schema.my_stream")

Elenco dei flussi

streams = client.list_streams(
    catalog_name="my_catalog",
    schema_name="my_schema",
    max_results=50,
    include_schemas=False,
)

Impostare include_schemas=True per includere i dettagli completi dello schema. Gli schemi possono essere di grandi dimensioni e ciò potrebbe comportare un'operazione a esecuzione prolungata. Per recuperare gli schemi singolarmente, usare get_streaminvece .

Eliminare un flusso

L'eliminazione di un flusso elimina anche la pipeline di acquisizione e la tabella di acquisizione.

Avvertimento

Tutti i modelli o le funzionalità che fanno riferimento al flusso eliminato non avranno più accesso ai dati del flusso sottostante. Creare una copia della tabella di inserimento prima dell'eliminazione se sono necessari questi dati, ma non è più necessario il flusso.

client.delete_stream(name="my_catalog.my_schema.my_stream")

Notebook di esempio

Per un esempio completo che crea uno Stream, definisce le funzionalità di streaming e distribuisce su un endpoint di serving, consulta il notebook seguente:

Notebook di avvio rapido sulle visualizzazioni delle funzionalità di streaming

Ottieni il notebook