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.
Important
Wersje środowiskowe dla potoków Lakeflow są dostępne w podglądzie publicznym.
Potoki z wersją environment set run Python code through Spark Connect. Ta strona omawia, co jest niekompatybilne, co zachowuje się inaczej oraz jak Databricks skanuje potok pod kątem dotkniętych wzorców.
Ograniczenia
Wersje środowiska nie są jeszcze zgodne ze wszystkimi funkcjami potoku. Uruchomienie potoku z zestawem wersji środowiska kończy się niepowodzeniem, jeśli kod Python potoku wykonuje dowolną z następujących czynności:
- Wycisza stan sesji platformy Spark wewnątrz funkcji ozdobionej dekoratorem potoków. Przykłady obejmują
spark.conf.set(...),spark.sql("USE CATALOG ...")icreateOrReplaceTempView. - Używa interfejsów API PySpark, które są niedostępne w programie Spark Connect, w tym
SparkContextinterfejsów API ,RDD,SQLContexti Py4J. Zobacz Co jest obsługiwane w programie Spark Connect.
Jeśli włączenie wersji środowiska w potoku powoduje jego niepowodzenie, wyłączenie wersji środowiska zwraca potok do poprzedniego stanu.
Zmiany zachowania
Program Spark Connect ma niewielką liczbę różnic w zachowaniu w klasycznym środowisku uruchomieniowym PySpark. Aby uzyskać pełną dokumentację, zobacz Spark Connect a klasyczna platforma Spark . Skan kompatybilności wykrywa te wzorce z wyprzedzeniem i blokuje migrację, dopóki nie zostaną rozwiązane, dzięki czemu można je znaleźć i naprawić, zanim wpłyną na dane produkcyjne.
W potoku najbardziej typowe sytuacje, w których zachowanie może się różnić, to:
- Przeplatane konstrukcje i mutacja sesji ramki danych
- UDFs odwołujące się do modyfikowalnego stanu Python
Przeplatane konstrukcje i mutacja sesji ramki danych
Gdy potok konstruuje ramkę danych, a następnie wycisza stan sesji platformy Spark (na przykład zmienia domyślny wykaz lub schemat, ustawia konfigurację, zastępuje widok tymczasowy lub ponownie rejestruje funkcję UDF), a następnie używa ramki danych:
- Bez wersji środowiska ramka danych używa stanu sesji wstępnej mutacji .
- W przypadku wersji środowiska ramka danych używa stanu sesji po mutacji .
Przykład:
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Bez wersji mytable środowiska zawiera [(1, "Original Row")]element . W przypadku wersji mytable środowiska zawiera element [(2, "Replaced Row")].
UDFs odwołujące się do stanu modyfikowalnego Python
Gdy funkcja zdefiniowanej przez użytkownika odwołuje się do zmiennej globalnej Python, której wartość zmienia się po zdefiniowaniu funkcji zdefiniowanej przez użytkownika:
- Bez wersji środowiska funkcja UDF używa najnowszej wartości zmiennej.
- W przypadku wersji środowiska funkcja UDF używa wartości w momencie, gdy zdefiniowano funkcję zdefiniowanej przez użytkownika.
Przykład:
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Bez wersji my_mv środowiska zawiera [("alex_b",)]element . W przypadku wersji my_mv środowiska zawiera element [("alex_a",)].
Jeśli potok korzysta z dowolnego wzorca, przeprowadź inspekcję przed włączeniem wersji środowiska.
Skanowanie zgodności
Skan zgodności znajduje wzorce kodu w twoim potoku, które dawałyby inne wyniki w wersji środowiskowej, więc możesz je naprawić przed automatyczną migracją potoku. Po włączeniu skanowania w potoku:
- Każda aktualizacja generuje jedno
BehaviorChangeInSparkConnectWARNzdarzenie w dzienniku zdarzeń potoku na każdy wykryty wzorzec. - Potok nie jest migrowany do wersji środowiskowej i nie możesz jej samodzielnie włączyć, dopóki wszystkie ostrzeżenia o kompatybilności z poprzedniej udanej aktualizacji nie zostaną usunięte.
Ta kontrola nie dotyczy potoku, który nie ma wcześniejszej aktualizacji ani który ma już ustawioną wersję środowiska.
Włączanie skanowania w potoku
Skanowanie zgodności można włączyć, dodając konfigurację potoku pipelines.environmentVersion.enableCompatibilityScan . Konfigurację można dodać za pomocą interfejsu użytkownika edytora potoku lub dodając wpis do pliku JSON konfiguracji potoku.
Za pośrednictwem interfejsu użytkownika:
- W edytorze rurociągu kliknij pozycję Ustawienia.
- Znajdź sekcję Konfiguracja w ustawieniach potoku.
- Kliknij
Dodaj konfigurację.
- Wprowadź
pipelines.environmentVersion.enableCompatibilityScanjako klucz itruejako wartość. - Zapisz ustawienia potoku.
W pliku JSON potoku:
Dodaj następujący wpis do configuration bloku:
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Przejrzyj i usuń ostrzeżenia dotyczące zgodności
Aby znaleźć i usunąć wzorce blokujące wersję środowiskową w Twoim potoku:
- Uruchom potok w trybie próby próbnej , a następnie zapytaj dziennik zdarzeń potoku o
BehaviorChangeInSparkConnectWARNzdarzenia. Każde zdarzenie zgłasza jeden wykryty wzorzec. Zobacz Dokumentację zdarzeń zgodności, aby zapoznać się z pełną listą kodów problemów, przykładowymi wzorcami i sugerowanymi poprawkami. - Zaktualizuj kod potoku, aby usunąć wykryte wzorce po sugerowanej poprawie, i uruchom potok ponownie.
- Powtarzaj, aż udana aktualizacja nie wygeneruje kolejnych zdarzeń kompatybilności. Potok może być następnie automatycznie migrowany, a także możesz samodzielnie włączyć wersję środowiskową.
Włączenie wersji środowiska wykonuje te same kontrole bezpieczeństwa, niezależnie od tego, czy Databricks automatycznie migruje potok, czy sam to ustawiasz.environment_version Potok z nierozwiązanymi ostrzeżeniami o zgodności nie przechodzi do wersji środowiskowej, dopóki ostrzeżenia nie zostaną rozwiązane. Jeśli migracja nie może zostać przeprowadzona bezpiecznie lub z jakiegokolwiek powodu zawiodła, system zatrzymuje się przed zapisem jakichkolwiek danych i potok kontynuuje działanie w poprzednim czasie działania.
Gdy aktualizacja zostaje zatrzymana z jednego z tych powodów, dziennik zdarzeń potoku oraz komunikat o błędzie aktualizacji opisują przyczynę oraz kroki jej rozwiązania. Postępuj zgodnie z tymi krokami i uruchom pipeline ponownie, aby zakończyć migrację. Jeśli uważasz, że ostrzeżenie o zgodności to fałszywy alarm, rozwiąż zaznaczony wzorzec lub skontaktuj się z pomocą wsparciem Azure Databricks.
Dokumentacja zdarzeń zgodności
Gdy skan zgodności uruchamia się na potoku, emituje jedno BehaviorChangeInSparkConnectWARN zdarzenie w dzienniku zdarzeń potoku na wykryty wzorzec. Gdy poprzednia udana aktualizacja wykryła jakiekolwiek wzorce, potok nie jest migrowany do wersji środowiskowej, dopóki wzorce nie zostaną rozwiązane.
Każde zdarzenie zgłasza pojedynczy kod problemu, który identyfikuje wykryte zdarzenia. Aby wyszukać kod, znajdź go w tabeli Kody problemów — każdy wiersz łączy się z sekcją kategorii zawierającą przykładowy wzorzec i sugerowaną poprawkę.
Kształt zdarzenia
BehaviorChangeInSparkConnect zdarzenia są zgodne ze standardowym schematem dziennika zdarzeń potoku:
- Parametr
event_typema wartośćbehavior_change_in_spark_connect. - Parametr
levelma wartośćWARN. -
detailsbehavior_change_in_spark_connectzawiera obiekt, który ma jednoissuepole. Wartość problemu to jeden z kodów wymienionych poniżej. -
messageto czytelny dla człowieka opis wykrytego wzorca.
Kody problemów
| Kategoria | Kod problemu | Description |
|---|---|---|
| Mutacje bazy danych i katalogu | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Domyślny wykaz został zmieniony po utworzeniu ramki danych. Istniejąca ramka danych może rozpoznawać tabele przy użyciu nowego wykazu domyślnego. |
| Mutacje bazy danych i katalogu | USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE CATALOG został wywołany poza funkcją ozdobioną przez dekorator potoków. Wykaz domyślny może zostać nieoczekiwanie zmieniony dla kolejnych operacji. |
| Mutacje bazy danych i katalogu | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Domyślna baza danych została zmieniona po utworzeniu ramki danych. Istniejąca ramka danych może rozpoznawać tabele przy użyciu nowej domyślnej bazy danych. |
| Mutacje bazy danych i katalogu | USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR |
USE DATABASE został wywołany poza funkcją ozdobioną przez dekorator potoków. Domyślna baza danych może zostać nieoczekiwanie zmieniona dla kolejnych operacji. |
| Chętne wykonywanie w ramach funkcji przepływu | CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wywołuje polecenie punktu kontrolnego. |
| Chętne wykonywanie w ramach funkcji przepływu | CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu chętnie tworzy widok ramki danych (createOrReplaceTempView lub podobny). |
| Chętne wykonywanie w ramach funkcji przepływu | CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu tworzy profil zasobu. |
| Chętne wykonywanie w ramach funkcji przepływu | GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wywołuje spark.resources lub powiązany interfejs API zasobów. |
| Chętne wykonywanie w ramach funkcji przepływu | MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wykonuje chętne MERGE INTO do tabeli docelowej. |
| Chętne wykonywanie w ramach funkcji przepływu | ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wykonuje chętną operację spark ML. |
| Chętne wykonywanie w ramach funkcji przepływu | REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu rejestruje źródło danych Python. |
| Chętne wykonywanie w ramach funkcji przepływu | STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu działa na aktywnym dojściu zapytania przesyłania strumieniowego. |
| Chętne wykonywanie w ramach funkcji przepływu | STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu rejestruje lub usuwa odbiornik zapytań przesyłania strumieniowego. |
| Chętne wykonywanie w ramach funkcji przepływu | STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wywołuje funkcję spark.streams do zarządzania zapytaniami przesyłanymi strumieniowo. |
| Chętne wykonywanie w ramach funkcji przepływu | WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wykonuje operację chętną DataFrameWriterV2 . |
| Chętne wykonywanie w ramach funkcji przepływu | WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu wykonuje operację chętną DataFrame.write . |
| Chętne wykonywanie w ramach funkcji przepływu | WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED |
Funkcja przepływu uruchamia zapytanie przesyłania strumieniowego (writeStream.start()). |
| Mutacje konfiguracji platformy Spark | CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED |
spark.conf.set() lub spark.conf.unset() został wywołany wewnątrz funkcji ozdobionej przez dekorator potoków. Nie jest to obsługiwane w wersji środowiska. |
| Mutacje konfiguracji platformy Spark | SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.set() została wywołana poza funkcją ozdobioną przez dekorator potoków po utworzeniu ramki danych. Zmiana konfiguracji może mieć wpływ na istniejącą ramkę danych w czasie wykonywania. |
| Mutacje konfiguracji platformy Spark | UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
spark.conf.unset() została wywołana poza funkcją ozdobioną przez dekorator potoków po utworzeniu ramki danych. Zmiana konfiguracji może mieć wpływ na istniejącą ramkę danych w czasie wykonywania. |
| Zamiany widoku tymczasowego | REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Globalny widok tymczasowy został zastąpiony po utworzeniu ramki danych odwołującej się do niej. Zastąpienie może zostać odzwierciedlone w istniejącej ramce danych. |
| Zamiany widoku tymczasowego | REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Widok tymczasowy został zastąpiony po utworzeniu ramki danych odwołującej się do niej. Zastąpienie może zostać odzwierciedlone w istniejącej ramce danych. |
| Mutacje UDF i UDTF | OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Funkcja UDF została ponownie zarejestrowana pod tą samą nazwą po utworzeniu ramki danych odwołującej się do niej. Istniejąca ramka danych może używać nowej definicji funkcji zdefiniowanej przez użytkownika. |
| Mutacje UDF i UDTF | OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR |
Funkcja UDTF została ponownie zarejestrowana pod tą samą nazwą po utworzeniu ramki danych odwołującej się do niej. Istniejąca ramka danych może używać nowej definicji udTF. |
| Mutacje UDF i UDTF | UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UDF odwołuje się do globalnej zmiennej modyfikowalnej Python. W przypadku wersji środowiska funkcja UDF używa wartości zmiennej w momencie, gdy zdefiniowano funkcję zdefiniowanej przez użytkownika, a nie w czasie wywołania. |
| Mutacje UDF i UDTF | UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR |
UTF odwołuje się do globalnej zmiennej Python modyfikowalnej. W przypadku wersji środowiska funkcja UDTF używa wartości zmiennej w momencie, gdy zdefiniowano zmienną UDTF, a nie w czasie wywołania. |
Mutacje bazy danych i katalogu
Te problemy są emitowane, gdy kod potoku wycisza domyślną bazę danych lub wykaz. W przypadku wersji środowiska ramki danych skonstruowane przed mutacją mogą rozpoznawać tabele przy użyciu nowej bazy danych lub wykazu.
Przykładowy wzorzec wyzwalający zdarzenie:
from pyspark import pipelines as dp
spark.sql("USE CATALOG marketing")
df = spark.read.table("events")
spark.sql("USE CATALOG sales") # changes the default catalog after df was created
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Bez wersji df środowiska jest rozpoznawana events z marketing katalogu. Wersja środowiska df jest rozpoznawana events z sales katalogu.
Sugerowana poprawka: W pełni kwalifikowane nazwy tabel, więc rozpoznawanie nie zależy od domyślnego wykazu lub bazy danych, i unikaj zmiany domyślnego wykazu lub bazy danych między tworzeniem i używaniem ramki danych.
from pyspark import pipelines as dp
df = spark.read.table("marketing.default.events")
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
Mutacje konfiguracji platformy Spark
Te problemy są emitowane, gdy kod potoku wycisza konfigurację platformy Spark w sposób, który może zmienić zachowanie ramki danych w wersji środowiska.
Przykładowy wzorzec wyzwalający zdarzenie:
from pyspark import pipelines as dp
df = spark.read.table("events")
spark.conf.set("spark.sql.ansi.enabled", "true") # changes session conf after df was created
@dp.materialized_view
def events_strict():
return df.selectExpr("CAST(price AS INT) AS price")
Bez wersji środowiska rzutowanie używa wartości conf w czasie tworzenia ramki danych. W przypadku wersji środowiska rzutowanie używa spark.sql.ansi.enabled=true metody i może zakończyć się niepowodzeniem w przypadku nieprawidłowych danych wejściowych.
Sugerowana poprawka: Ustaw wszystkie wymagane konfiguracje platformy Spark w górnej części pliku potoku przed utworzeniem dowolnej ramki danych. W przypadku konfiguracji poszczególnych zapytań użyj ustawienia potoku configuration w specyfikacji potoku.
Zamiany widoku tymczasowego
Te problemy są emitowane, gdy kod potoku zastępuje tymczasowy widok po utworzeniu ramki danych odwołującej się do niego. W przypadku wersji środowiska istniejąca ramka danych może odzwierciedlać nową zawartość widoku.
Przykładowy wzorzec wyzwalający zdarzenie:
from pyspark import pipelines as dp
spark.createDataFrame([(1, "Original Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
df = spark.sql("SELECT * FROM my_view")
spark.createDataFrame([(2, "Replaced Row")], ["id", "data"]) \
.createOrReplaceTempView("my_view")
@dp.materialized_view
def mytable():
return df
Bez wersji mytable środowiska zawiera [(1, "Original Row")]element . W przypadku wersji mytable środowiska zawiera element [(2, "Replaced Row")].
Sugerowana poprawka: Utwórz każdy widok tymczasowy pojedynczo i nie zastąp go. Jeśli potrzebujesz wielu widoków z powiązanymi danymi, nadaj każdej innej nazwie.
Mutacje UDF i UDTF
Te problemy są emitowane, gdy kod potoku wycisza funkcję UDF lub UDTF w sposób, który zmienia zachowanie w wersji środowiska.
Przykładowy wzorzec wyzwalający zdarzenie:
from pyspark import pipelines as dp
from pyspark.sql.functions import col, udf
suffix = "a"
@udf
def my_udf(s):
return s + suffix
suffix = "b"
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(my_udf(col("name")))
Bez wersji my_mv środowiska zawiera [("alex_b",)]element . W przypadku wersji my_mv środowiska zawiera element [("alex_a",)].
Suggested fix: Przekaż wartości do funkcji zdefiniowanej przez użytkownika jako argumenty zamiast przechwytywania ich z Python globalnych lub ustaw globalne przed zdefiniowaniem funkcji zdefiniowanej przez użytkownika i nie zmutuj ich później.
from pyspark import pipelines as dp
from pyspark.sql.functions import col, lit, udf
@udf
def append_suffix(s, suffix):
return s + suffix
@dp.materialized_view
def my_mv():
return spark.createDataFrame([("alex",)], ["name"]).select(append_suffix(col("name"), lit("b")))
Chętne wykonywanie w ramach funkcji przepływu
Te problemy są emitowane, gdy kod potoku wykonuje chętne polecenie spark wewnątrz funkcji ozdobionej przez dekorator potoków (@table, @materialized_viewitp.). Oczekuje się, że funkcje przepływu zdefiniują i zwracają ramkę danych; chętne polecenia, które zapisują dane, zarządzają zapytaniami przesyłanymi strumieniowo, rejestrują zasoby lub uruchamiają operacje uczenia maszynowego, nie są dozwolone w funkcji przepływu z zestawem wersji środowiska.
Sugerowana poprawka: Przenieś operację chętną poza funkcję przepływu i zamiast tego zwróć ramkę danych z funkcji przepływu. Skutki uboczne, takie jak zapisywanie w tabeli lub uruchamianie zapytania przesyłania strumieniowego, należą do definicji potoku; aparat potoku obsługuje materializację ramki danych zwracanej przez funkcję przepływu.
Znajdowanie zdarzeń zgodności w dzienniku zdarzeń
Następujące zapytanie zwraca wszystkie zdarzenia zgodności dla potoku uporządkowane po raz pierwszy:
SELECT
timestamp,
message,
details:behavior_change_in_spark_connect:issue AS issue
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
ORDER BY timestamp DESC;
Aby zliczyć zdarzenia według kodu problemu w ostatnich aktualizacjach:
SELECT
details:behavior_change_in_spark_connect:issue AS issue,
COUNT(*) AS occurrences
FROM event_log(<pipeline-id>)
WHERE event_type = 'behavior_change_in_spark_connect'
AND level = 'WARN'
GROUP BY 1
ORDER BY occurrences DESC;
Aby dowiedzieć się, jak wykonywać zapytania dotyczące dziennika zdarzeń, zobacz Wykonywanie zapytań w dzienniku zdarzeń.
Dodatkowe zasoby
- Konfiguruj wersje środowisk dla potoków — przegląd funkcji, automatyczna migracja oraz sposób samodzielnego włączenia wersji środowiska.
- Schemat dziennika zdarzeń potoku — schemat dziennika zdarzeń pełnego potoku.
- Dziennik zdarzeń potoku — jak wykonywać zapytania dotyczące dziennika zdarzeń potoku.