Nawiązywanie połączenia z usługą Lakebase

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. Element catalog to 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ład id lub user_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 format projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Zobacz Identyfikatory obliczeniowe.
  • <database>:Fakultatywny. Nazwa docelowej bazy danych PostgreSQL. Wartość domyślna to databricks_postgres. Zobacz Zarządzanie bazami danych.
  • <schema>.<table>: tabela docelowa w schema.table formacie. Jeśli pominięto schemat, ujście używa schematu public . 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ład id lub user_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 w schema.table formacie. Jeśli pominięto schemat, ujście używa schematu public .
  • <primary-key-columns>:Fakultatywny. Lista oddzielona przecinkami wszystkich kolumn w kluczu głównym tabeli docelowej, na przykład id lub user_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 korzysta sslmode=verify-full z 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 batchsize wierszy, a wartość domyślna to 1000.
  • Wiek buforu przekracza batchintervalwartość , która jest domyślnie ustawiona na 100 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.