Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
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-clientpacchetto 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
Connessione al catalogo Unity (scelta consigliata)
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:
-
SELECTnella tabella di acquisizione concede l'accesso in lettura allo stream. -
MANAGEnella 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_timestampnell'origine del riempimento. -
Inserimento min: numero minimo
stream_record_timestampdi 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_partitionekafka_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,valueestream_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: