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 Structured Streaming, aby zapisywać dane do Lakebase lub zewnętrznej bazy danych PostgreSQL z wbudowanym przetwarzaniem wsadowym, automatycznym ponawianiem prób i uwierzytelnianiem zarządzanym przez obszar roboczy.
Kiedy należy użyć ujścia Lakebase
Użyj ujścia Lakebase do strumieniowego zapisu danych z niskimi opóźnieniami do Lakebase lub zewnętrznej bazy danych PostgreSQL. Ten odbiornik nie wymaga implementowania niestandardowych funkcji foreach do obsługi przetwarzania wsadowego, zarządzania połączeniami i obsługi błędów.
Typowe przypadki użycia to:
- Aktualizowanie baz danych aplikacji w czasie rzeczywistym na potrzeby operacyjnych pulpitów nawigacyjnych lub funkcji dostępnych dla klientów.
- Synchronizuj stale zmieniające się dane, takie jak zagregowane lub filtrowane wyniki przesyłania strumieniowego, do transakcyjnej bazy danych.
- Zapisz wyniki zapytania Structured Streaming w tabeli Lakebase z opóźnieniem poniżej sekundy przy użyciu trybu czasu rzeczywistego.
Aby synchronizować dane z Lakebase do tabel Delta Lake w usłudze Lakehouse, w przeciwnym kierunku zobacz Źródło danych o zmianach w Lakebase.
Wymagania
-
Databricks Runtime 18 LTS i wyższe.
- Zewnętrzne połączenia PostgreSQL wymagają korzystania z Databricks Runtime 19 w wersji 19 lub nowszej oraz zgłoszenia się do wersji zapoznawczej Custom JDBC on UC Compute.
- Typy danych interwałowych wymagają korzystania z Databricks Runtime 19 lub nowszej wersji.
- Tradycyjne zasoby obliczeniowe z dedykowanym lub standardowym trybem dostępu albo obliczenia bezserwerowe dla notebooków lub zadań. W przypadku przetwarzania bezserwerowego użyj
Trigger.AvailableNow(). Zobacz Przesyłanie strumieniowe w obliczeniach bezserwerowych. - Baza danych Lakebase lub połączenie Unity Catalog z zewnętrzną bazą PostgreSQL.
Wymagania dotyczące identyfikatorów
Dla wszystkich celów Databricks zaleca stosowanie nazw kolumn typu schemat, tabela, kolumn oraz klucza głównego, które zaczynają się na literę lub podkreślenie i zawierają jedynie litery, cyfry i podkreślenia. Odbiornik wymusza spełnienie tych wymagań podczas automatycznego tworzenia tabeli Lakebase. Aby użyć identyfikatorów, które nie spełniają tych wymagań, utwórz tabelę docelową przed rozpoczęciem zapytania.
Połącz z bazą danych
Ujście usługi Lakebase obsługuje następujące metody połączenia:
Tabele Lakebase zarejestrowane w Unity Catalog
W przypadku tabel Lakebase zarejestrowanych w usłudze Unity Catalog łącznik automatycznie zarządza poświadczeniami i posługuje się tożsamością użytkownika lub nazwy głównej usługi, która uruchamia zapytanie. Jeśli tabela nie istnieje, łącznik tworzy tabelę.
Aby zarejestrować bazę danych Lakebase w Unity Catalog, zobacz Rejestrowanie bazy danych Lakebase w Unity Catalog.
Aby zapisywać do tabeli Lakebase, użyj metody .toTable() i podaj w pełni kwalifikowaną nazwę tabeli: catalog.schema.table
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
Zastąp następujące symbole zastępcze:
-
<catalog>.<schema>.<table>: To w pełni kwalifikowana nazwa tabeli docelowej. Elementcatalogto katalog w Unity Catalog utworzony podczas rejestrowania bazy danych Lakebase, zobacz Rejestrowanie bazy danych Lakebase w Unity Catalog. Jeśli tabela nie istnieje, łącznik go utworzy. -
<primary-key-columns>:Fakultatywny. Lista oddzielona przecinkami wszystkich kolumn w kluczu głównym tabeli docelowej, na przykładidlubuser_id,event_type. Zobacz działanie funkcji Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Ścieżka wolumenu Unity Catalog, w której zapytanie zapisuje swój punkt kontrolny. Możesz również użyć adresu URI magazynu obiektów w chmurze. Lokalizacja musi być magazynem, do którego można zapisywać dane, a nie na dysku lokalnym i musi być unikatowa dla każdego zapytania przesyłania strumieniowego. Jest to niezależne od tabeli docelowej. Zobacz Ustrukturyzowane punkty kontrolne przesyłania strumieniowego.
Aby uzyskać informacje o opcjonalnych konfiguracjach, takich jak batchsize i batchinterval, zobacz opcje ujścia PostgreSQL.
Tabele Lakebase nie zostały zarejestrowane w Unity Catalog
W przypadku tabel Lakebase, które nie są zarejestrowane w Unity Catalog, łącznik automatycznie zarządza poświadczeniami i korzysta z tożsamości użytkownika lub nazwy głównej usługi, która uruchamia zapytanie. Jeśli tabela nie istnieje, łącznik tworzy tabelę.
Aby zapisać w tabeli Lakebase, użyj opcji endpoint i dbtable:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") // Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
Zastąp następujące symbole zastępcze:
-
<project-id>.<branch-id>.<endpoint-id>: Punkt końcowy usługi Lakebase. Znajdź wszystkie trzy wartości w pozycji Nazwa zasobu w menu Pobierz identyfikator na karcie Computes, która ma formatprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Zobacz Identyfikatory obliczeniowe. -
<database>:Fakultatywny. Nazwa docelowej bazy danych PostgreSQL. Wartość domyślna todatabricks_postgres. Zobacz Zarządzanie bazami danych. -
<schema>.<table>: tabela docelowa wschema.tableformacie. Jeśli pominięto schemat, ujście używa schematupublic. Do automatycznego tworzenia tabel używaj identyfikatorów zaczynających się od litery lub podkreślenia, zawierających wyłącznie litery, cyfry i podkreślenia. -
<primary-key-columns>:Fakultatywny. Lista oddzielona przecinkami wszystkich kolumn w kluczu głównym tabeli docelowej, na przykładidlubuser_id,event_type. Zobacz działanie funkcji Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Ścieżka wolumenu Unity Catalog, w której zapytanie zapisuje swój punkt kontrolny. Możesz również użyć adresu URI magazynu obiektów w chmurze. Lokalizacja musi być magazynem, do którego można zapisywać dane, a nie na dysku lokalnym i musi być unikatowa dla każdego zapytania przesyłania strumieniowego. Jest to niezależne od tabeli docelowej. Zobacz Ustrukturyzowane punkty kontrolne przesyłania strumieniowego.
Informacje o opcjonalnych konfiguracjach, takich jak batchsize i batchinterval, można znaleźć w sekcji opcje ujścia PostgreSQL.
Zewnętrzny PostgreSQL z poświadczeniami Unity Catalog
Ważna
Ta funkcja jest dostępna w publicznej wersji testowej. Administratorzy obszaru roboczego mogą zarządzać dostępem do niestandardowego JDBC w UC Compute na stronie Wersje zapoznawcze. Zobacz Zarządzanie wersjami zapoznawczami usługi Azure Databricks.
Użyj połączenia Unity Catalog, aby uwierzytelnić się do zewnętrznej bazy PostgreSQL bez przechowywania danych w kodzie. Tabela docelowa musi już istnieć.
Stwórz połączenie typu POSTGRESQL, zobacz Utwórz połączenie. Użytkownik lub główny użytkownik usługi wykonujący zapytanie musi mieć USE CONNECTION na połączeniu.
Aby zapisywać dane w tabeli PostgreSQL, użyj opcji databricks.connection, database i dbtable:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
Zastąp następujące symbole zastępcze:
-
<connection-name>: Nazwa połączenia z katalogiem Unity. -
<database>: Nazwa docelowej bazy danych PostgreSQL. -
<schema>.<table>: Istniejąca tabela docelowa wschema.tableformacie. Jeśli pominięto schemat, ujście używa schematupublic. -
<primary-key-columns>:Fakultatywny. Lista oddzielona przecinkami wszystkich kolumn w kluczu głównym tabeli docelowej, na przykładidlubuser_id,event_type. Zobacz działanie funkcji Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Ścieżka wolumenu Unity Catalog, w której zapytanie zapisuje swój punkt kontrolny. Możesz również użyć adresu URI magazynu obiektów w chmurze. Lokalizacja musi być magazynem, do którego można zapisywać dane, a nie na dysku lokalnym i musi być unikatowa dla każdego zapytania przesyłania strumieniowego. Jest to niezależne od tabeli docelowej. Zobacz Ustrukturyzowane punkty kontrolne przesyłania strumieniowego.
Połączenia PostgreSQL zawsze korzystają z TLS. Weryfikacja certyfikatu przebiega zgodnie z ustawieniami połączenia z Unity Catalog, które wybierasz podczas tworzenia połączenia:
-
Zaufaj certyfikatowi serwera: Po wybraniu połączenie używa
sslmode=require, co szyfruje połączenie bez weryfikacji certyfikatu serwera. -
Certyfikat serwera dostarczony przez użytkownika: Udostępnij certyfikat serwera zakodowany PEM, który można użyć
sslmode=verify-full, gdy certyfikat serwera Trust nie jest wybrany. Jeśli nie podasz certyfikatu, połączenie korzystasslmode=verify-fullz domyślnego magazynu zaufania JVM.
Opcje konfiguracji
Odbiornik zgłasza błąd w przypadku nierozpoznanych opcji, JDBC_STREAMING_SINK_INVALID_OPTIONS.
Aby poznać opcje konfiguracji sinka, w tym wspólne opcje oraz opcje dla każdej metody połączenia, zobacz opcje sinka PostgreSQL.
Mapowanie typu danych
Sink sprawdza, czy każda kolumna DataFrame jest zgodna z odpowiadającą jej kolumną docelową, zanim zapisze ją do istniejącej tabeli Lakebase lub zewnętrznej PostgreSQL.
Poniższa tabela zawiera typy obsługiwane w Databricks Runtime 18 LTS i wyższych:
| Typ Spark | Automatycznie tworzony typ tabeli Lakebase | Kompatybilne typy w istniejących tabelach PostgreSQL |
|---|---|---|
ByteType, ShortType |
smallint |
smallint |
IntegerType |
integer |
integer |
LongType |
bigint |
bigint |
FloatType |
real |
real |
DoubleType |
double precision |
double precision |
DecimalType |
numeric |
numeric |
StringType |
text |
varchar, text |
VarcharType(n) |
varchar(n) |
varchar, text |
CharType(n) |
char(n) |
char |
BinaryType |
bytea |
bytea |
BooleanType |
boolean |
boolean |
TimestampType |
timestamptz |
timestamptz |
TimestampNTZType |
timestamp |
timestamp |
DateType |
date |
date |
ArrayType, MapType, , StructType, VariantTypeNullType |
jsonb |
json, jsonb |
Poniższa tabela zawiera typy obsługiwane w Databricks Runtime 19 i wyższych:
| Typ Spark | Automatycznie tworzony typ tabeli Lakebase | Kompatybilne typy w istniejących tabelach PostgreSQL |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Zachowanie operacji upsert
Opcja upsertkey ta identyfikuje główne kolumny klucza tabeli docelowej. Dla istniejącej tabeli kolumny w upsertkey muszą dokładnie odpowiadać kluczowi podstawowemu tabeli. Jeśli pominiesz tę opcję, zlew odczytuje klucz główny ze stołu. W przypadku tabeli Lakebase tworzonej przez sink element upsertkey definiuje klucz główny. Jeśli pominiesz tę opcję, ujście utworzy tabelę bez klucza podstawowego.
Gdy tabela docelowa ma klucz główny, ujście wykonuje operacje upsert przy użyciu składni PostgreSQL INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET .... Gdy tabela docelowa nie ma klucza głównego, sink wykonuje wstawki. Tryb wyjścia zapytania nie ma wpływu na to zachowanie.
Wszystkie kolumny klucza podstawowego muszą być obecne w obiekcie DataFrame i mieć porównywalne typy, takie jak typy liczbowe lub ciągi znaków.
Optymalizacja wydajności
Przetwarzanie wsadowe i mechanizm przeciwprzeciążeniowy
Opróżnianie jest wyzwalane po spełnieniu dowolnego warunku:
- Bufor osiąga
batchsizewierszy, a wartość domyślna to1000. - Wiek buforu przekracza
batchintervalwartość , która jest domyślnie ustawiona na100 milliseconds.
Gdy baza danych nie może nadążyć za tempem napływu danych, odbiornik propaguje presję wsteczną w górę strumienia do źródła.
Wskazówki dotyczące opóźnienia i przepływności:
- W przypadku obciążeń o małych opóźnieniach w trybie czasu rzeczywistego zmniejsz rozmiar
batchinterval, aby zagwarantować krótszy maksymalny czas przed opróżniniem. Zobacz koncepcje trybu czasu rzeczywistego dla koncepcji oraz przykłady trybu czasu rzeczywistego dla przykładu kodu. - W przypadku obciążeń o dużej przepustowości zwiększ
batchsize, aby zmniejszyć narzut dla każdej transakcji.
Zachowanie połączenia
Ujście używa buforowania połączeń w funkcjach wykonawczych. Domyślnie każde zadanie używa jednego połączenia z bazą danych.
Databricks zaleca używanie domyślnej wartości zadania 1 dla każdego połączenia. Jeśli zwiększysz liczbę zadań dla każdego połączenia, może to spowodować rywalizację o połączenia i zwiększyć opóźnienia dla połączeń o wysokiej przepływności.
Aby skonfigurować stosunek zadań do połączeń, ustaw konfigurację platformy spark.databricks.sql.streaming.jdbc.tasksPerConnection Spark. Jeśli docelowa baza danych ma niski limit połączenia, zmniejsz liczbę partycji mieszania lub zwiększ wartość spark.databricks.sql.streaming.jdbc.tasksPerConnection.
Sink automatycznie ponawia próby w przypadku przejściowych błędów JDBC, w tym błędów połączenia, zakleszczeń i ograniczeń szybkości. Jeśli ujście wyczerpuje wszystkie ponawianie prób, zapytanie zakończy się niepowodzeniem.
Obsługiwane wyzwalacze i tryby wyjściowe
Triggers
Ta tabela przedstawia obsługę typów wyzwalaczy w Structured Streaming w klasycznych i bezserwerowych zasobach obliczeniowych:
| Trigger | Obliczenia klasyczne | Bezserwerowe przetwarzanie (notebooki i zadania) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Yes. Deprecated. Użyj AvailableNow. |
Yes. Deprecated. Użyj AvailableNow. |
Tryby wyjściowe
Ta tabela przedstawia obsługę trybów wyjścia funkcji Structured Streaming:
| Tryb wyjściowy | Supported |
|---|---|
update |
Yes |
append |
Yes. Zachowanie jest identyczne z update. Zapytanie wykonuje operację upsert, gdy tabela docelowa ma klucz podstawowy; w przeciwnym razie wykonuje operację wstawiania. Zobacz działanie funkcji Upsert. |
complete |
No |
Ograniczenia
- Dla zewnętrznej bazy danych PostgreSQL połączonej przez połączenie Unity Catalog, docelowa tabela musi już istnieć. Zlew automatycznie tworzy brakujące tabele tylko w Lakebase.
- Rurociągi przepływowe do jezior nie są obsługiwane.