Notatka
Dostęp do tej strony wymaga autoryzacji. Może spróbować zalogować się lub zmienić katalogi.
Dostęp do tej strony wymaga autoryzacji. Możesz spróbować zmienić katalogi.
Użyj interfejsu API potoku Lakeflow sink z przepływami, aby zapisywać rekordy przekształcone przez potok w zewnętrznym miejscu docelowym danych. Zewnętrzne źródła danych obejmują zarządzane przez Unity Catalog i zewnętrzne tabele oraz usługi przesyłania strumieniowego zdarzeń, takie jak Apache Kafka lub Azure Event Hubs. Za pomocą pochłaniaczy danych można również zapisywać do niestandardowych źródeł danych, pisząc kod w języku Python dla tych źródeł danych.
Aby zapoznać się z omówieniem pojęć ujścia i kiedy ich używać, zobacz Sinks in Lakeflow pipelines (Ujścia w potokach lakeflow).
Uwaga / Notatka
- Interfejs
sinkAPI jest dostępny tylko dla języka Python. - Możesz utworzyć niestandardowy "sink" przy użyciu interfejsu API ForEachBatch. Zobacz Use ForEachBatch to write to arbitrary data sinks in pipelines (Używanie funkcji ForEachBatch do zapisywania w dowolnych ujściach danych w potokach).
Przepływ pracy ujścia danych
Gdy dane zdarzenia są pozyskiwane ze źródła przesyłania strumieniowego do potoku, przetwarzasz i uściślisz te dane w przekształceniach w potoku. Następnie używasz przetwarzania przepływu typu "append", aby przesyłać strumieniowo przekształcone rekordy danych do miejsca docelowego. Tworzysz ten zlew przy użyciu funkcji create_sink(). Aby uzyskać więcej informacji na temat funkcji, zobacz dokumentację interfejsu create_sinkAPI ujścia.
Jeśli masz potok danych, który tworzy lub przetwarza dane zdarzeń przesyłania strumieniowego i przygotowuje rekordy danych do zapisu, możesz użyć ujścia.
Implementowanie ujścia składa się z dwóch kroków:
- Utwórz ujście.
- Użyj przepływu dołączania lub przepływu aktualizacji , aby zapisać przygotowane rekordy do ujścia.
Tworzenie ujścia
Usługa Databricks obsługuje kilka typów ujściów docelowych, w których zapisujesz rekordy przetwarzane z danych strumienia:
- Miejsca docelowe dla tabel Delta (w tym tabele zarządzane przez Unity Catalog i tabele zewnętrzne)
- Ujścia platformy Apache Kafka
- Ujścia usługi Azure Event Hubs
- Niestandardowe ujścia napisane w języku Python przy użyciu niestandardowych źródeł danych języka Python
Poniżej przedstawiono przykłady konfiguracji dla ujścia usług Delta, Kafka i Azure Event Hubs oraz niestandardowych źródeł danych języka Python:
Rozlewiska delty
Aby utworzyć sink Delta za pomocą ścieżki pliku:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = {"path": "/Volumes/catalog_name/schema_name/volume_name/path/to/data"}
)
Aby utworzyć Delta sink na podstawie nazwy tabeli przy użyciu w pełni kwalifikowanej ścieżki katalogu i schematu:
dp.create_sink(
name = "delta_sink",
format = "delta",
options = { "tableName": "catalog_name.schema_name.table_name" }
)
Ujścia platformy Kafka i usługi Azure Event Hubs
Ten kod działa zarówno dla ujść Apache Kafka, jak i Azure Event Hubs.
credential_name = "<service-credential>"
eh_namespace_name = "dp-eventhub"
bootstrap_servers = f"{eh_namespace_name}.servicebus.windows.net:9093"
topic_name = "dp-sink"
dp.create_sink(
name = "eh_sink",
format = "kafka",
options = {
"databricks.serviceCredential": credential_name,
"kafka.bootstrap.servers": bootstrap_servers,
"topic": topic_name
}
)
To jest credential_name, czyli odwołanie do poświadczenia usługi katalogu Unity. Aby uzyskać więcej informacji, patrz Używanie poświadczeń usługi Unity Catalog do łączenia się z zewnętrznymi usługami w chmurze.
Niestandardowe źródła danych języka Python
Zakładając, że masz niestandardowe źródło danych języka Python zarejestrowane jako my_custom_datasource, poniższy kod może zapisywać dane w tym źródle danych.
from pyspark import pipelines as dp
# Assume `my_custom_datasource` is a custom Python streaming
# data source that writes data to your system.
# Create Lakeflow pipelines sink using my_custom_datasource
dp.create_sink(
name="custom_sink",
format="my_custom_datasource",
options={
<options-needed-for-custom-datasource>
}
)
# Create append flow to send data to RequestBin
@dp.append_flow(name="flow_to_custom_sink", target="custom_sink")
def flow_to_custom_sink():
return read_stream("my_source_data")
Aby uzyskać szczegółowe informacje na temat tworzenia niestandardowych źródeł danych w języku Python, zobacz PySpark custom data sources (Niestandardowe źródła danych PySpark).
Aby uzyskać więcej informacji na temat korzystania z funkcji create_sink, zobacz dokumentację interfejsu API sink .
Po utworzeniu ujścia można rozpocząć przesyłanie przetworzonych rekordów do niego strumieniowo.
Zapisz do zbiornika przy użyciu przepływu dołączania
Po utworzeniu ujścia, następnym krokiem jest zapisanie przetworzonych rekordów, wskazując je jako miejsce docelowe dla rekordów wyprowadzanych przez przepływ dodawania. Aby to zrobić, określasz zlewozmywak jako wartość target w dekoratorze append_flow.
- W przypadku tabel zarządzanych i zewnętrznych Unity Catalog użyj formatu
deltai w opcjach określ ścieżkę lub nazwę tabeli. Potok danych musi być skonfigurowany do korzystania z Unity Catalogu. - W przypadku tematów Apache Kafka użyj formatu
kafkai określ w opcjach nazwę tematu, informacje o połączeniu oraz informacje uwierzytelniające. Są to te same opcje obsługiwane przez ujście Kafka dla Spark Structured Streaming. Zobacz Konfigurowanie zapisywania przesyłania strumieniowego w strukturze na platformie Kafka. - W przypadku usługi Azure Event Hubs użyj formatu
kafkai określ nazwę usługi Event Hubs, informacje o połączeniu i informacje dotyczące uwierzytelniania w opcjach. Są to te same opcje obsługiwane w ujściu strumieniowania strukturalnego Spark Event Hubs, które korzysta z interfejsu Kafka. Zobacz Uwierzytelnianie.
Poniżej przedstawiono przykłady konfigurowania przepływów do zapisywania w ujściach usługi Delta, Kafka i Azure Event Hubs z rekordami przetwarzanymi przez potok.
Odbiornik delta
@dp.append_flow(name = "delta_sink_flow", target="delta_sink")
def delta_sink_flow():
return(
spark.readStream.table("spark_referrers")
.selectExpr("current_page_id", "referrer", "current_page_title", "click_count")
)
Ujścia platformy Kafka i usługi Azure Event Hubs
@dp.append_flow(name = "kafka_sink_flow", target = "eh_sink")
def kafka_sink_flow():
return (
spark.readStream.table("spark_referrers")
.selectExpr("cast(current_page_id as string) as key", "to_json(struct(referrer, current_page_title, click_count)) AS value")
)
Parametr value jest obowiązkowy dla ujścia usługi Azure Event Hubs. Dodatkowe parametry, takie jak key, partition, headersi topic, są opcjonalne.
Aby uzyskać więcej informacji na temat dekoratora append_flow, zobacz Przepływy domyślne i dołączane.
Ograniczenia
Obsługiwany jest tylko interfejs API języka Python. Język SQL nie jest obsługiwany.
Obsługiwane są tylko zapytania przesyłane strumieniowo. Zapytania wsadowe nie są obsługiwane.
Tylko
append_flowiupdate_flowmogą być używane do zapisu do sinków. Inne przepływy, takie jakcreate_auto_cdc_flow, nie są obsługiwane i nie można użyć ujścia w definicji zestawu danych potoku. Na przykład następujące nie są obsługiwane:@table("from_sink_table") def fromSink(): return read_stream("my_sink")W przypadku systemów Delta nazwa tabeli musi być w pełni kwalifikowana. W szczególności w przypadku tabel zewnętrznych zarządzanych przez Unity Catalog, nazwa tabeli musi mieć format
<catalog>.<schema>.<table>. W przypadku magazynu metadanych Hive musi on znajdować się w postaci<schema>.<table>.Uruchomienie aktualizacji pełnego odświeżania nie powoduje wyczyszczenia wcześniej obliczonych danych wyników w ujściach. Oznacza to, że wszystkie ponownie przetworzone dane są dołączane do ujścia, a istniejące dane nie są zmieniane.
Oczekiwania dotyczące kanału przetwarzania nie są obsługiwane.
Bezserwerowa kontrola ruchu wychodzącego obsługuje tylko konektory docelowe Kafka i Delta Lake. Zobacz Co to jest kontrola ruchu wychodzącego w środowisku bezserwerowym?.
Dodatkowe zasoby
- Potoki deklaratywne platformy Spark
- Co to są potoki?
- Ładuj i przetwarzaj dane stopniowo za pomocą przepływów potoku Lakeflow
- dokumentacja interfejsu API ujścia języka Python