Używanie Unity Catalog z przesyłaniem strumieniowym o strukturze

Na tej stronie opisano, jak używać Structured Streaming z Unity Catalog do zarządzania nadzorem nad danymi w przypadku obciążeń przyrostowych i strumieniowych w usłudze Azure Databricks.

Jaką funkcjonalność Structured Streaming obsługuje Unity Catalog?

Unity Catalog nie wprowadza żadnych wyraźnych limitów dla źródeł i ujść Structured Streaming dostępnych w Azure Databricks.

Dzięki Unity Catalog i Structured Streaming możesz:

  • Przesyłaj strumieniowo dane zarówno z tabel zarządzanych, jak i zewnętrznych. Zobacz tabele zarządzane w Unity Catalog dla Delta Lake i Apache Iceberg.
  • Korzystaj z lokalizacji zewnętrznych zarządzanych przez Unity Catalog do interakcji z danymi przy użyciu identyfikatorów URI magazynu obiektów.
  • Zapisuj do tabel zewnętrznych, używając nazw tabel albo ścieżek plików. Aby wchodzić w interakcje z tabelami zarządzanymi, musisz użyć nazwy tabeli.

W przypadku punktów kontrolnych Structured Streaming należy używać ścieżek w lokalizacjach zewnętrznych zarządzanych przez Unity Catalog. Aby dowiedzieć się więcej na temat bezpiecznego łączenia magazynu z Unity Catalog, zobacz część Połącz się z magazynem obiektów w chmurze przy użyciu Unity Catalog.

odczytać widok Unity Catalog jako strumień

W Databricks Runtime 14.3 LTS i nowszych wersjach można używać Structured Streaming do odczytywania danych z widoków zarejestrowanych w Unity Catalog. Tabele źródłowe muszą korzystać z formatu Delta Lake. Aby uzyskać informacje o innych ograniczeniach, zobacz Ograniczenia.

Aby odczytać widok przy użyciu Structured Streaming, użyj metody .table() za pomocą identyfikatora widoku:

df = (spark.readStream
  .table("demoView")
)

Użytkownicy muszą mieć SELECT uprawnienia w widoku docelowym.

Jeśli zmodyfikujesz definicję widoku, aby dodać lub zmienić tabele, do których odwołujesz się w widoku, nie możesz użyć tego samego punktu kontrolnego przesyłania strumieniowego.

Obsługiwane opcje przesyłania strumieniowego

Czytnik strumieniowy stosuje opcje do plików i metadanych bazowych tabel Delta Lake dla określonego widoku.

Obsługiwane są następujące opcje:

  • maxFilesPerTrigger
  • maxBytesPerTrigger
  • ignoreDeletes
  • skipChangeCommits
  • withEventTimeOrder
  • startingTimestamp
  • startingVersion

Odczyty w widokach z UNION ALL nie obsługują opcji withEventTimeOrder i startingVersion.

Jeśli podasz nieobsługiwane opcje, takie jak readChangeFeed, platforma Spark zgłosi ten wyjątek:

AnalysisException: [UNSUPPORTED_STREAMING_OPTIONS_FOR_VIEW.UNSUPPORTED_OPTION] Unsupported for streaming a view. Reason: option <option> is not supported.

Obsługiwane operacje przesyłania strumieniowego

Obsługiwane operacje obejmują:

Operation Description Operator Example
Projekt Kontroluje uprawnienia na poziomie kolumny SELECT... FROM... CREATE VIEW project_view AS SELECT id, value FROM source_table
Filtr Steruje uprawnieniami na poziomie wiersza WHERE... CREATE VIEW filter_view AS SELECT * FROM source_table WHERE value > 100
Wszystkie unii Wyniki z wielu tabel UNION ALL CREATE VIEW union_view AS SELECT id, value FROM source_table1 UNION ALL SELECT * FROM source_table2

Nieobsługiwane operacje obejmują agregacje, sortowanie i funkcje zwracające tabele, takie jak table_changes(). Szczegółowe informacje na temat funkcji zwracających tabelę można znaleźć w sekcji Wywołanie funkcji zwracającej tabelę (TVF).

Jeśli przetwarzasz strumieniowo dane z widoku zawierającego nieobsługiwaną operację, Spark zgłasza następujący wyjątek:

UnsupportedOperationException: [UNEXPECTED_OPERATOR_IN_STREAMING_VIEW] Unexpected operator <operator> in the CREATE VIEW statement as a streaming source. A streaming view query must consist only of SELECT, WHERE, and UNION ALL operations.

Ograniczenia