Compatibilidade de versão do ambiente
Pré-lançamento público
As versões de ambiente para LakeFlow Pipelines estão em pré-lançamento público.
Pipelines com uma versão de ambiente definida executam código Python por meio do Spark Connect. Esta página aborda o que é incompatível, o que se comporta de maneira diferente e como o Databricks verifica um pipeline em busca de padrões afetados.
Limitações
As versões de ambiente ainda não são compatíveis com todas as funcionalidades do pipeline. A execução de um pipeline com uma versão de ambiente definida falha se o código Python do pipeline fizer qualquer uma das seguintes ações:
- Altera o estado da sessão Spark dentro de uma função decorada com um decorador de pipeline. Exemplos incluem
spark.conf.set(...),spark.sql("USE CATALOG ...")ecreateOrReplaceTempView. - Usa APIs PySpark que não estão disponíveis no Spark Connect, incluindo
SparkContext,RDD,SQLContexte quaisquer APIs Py4J. Consulte a seção "O que é compatível com o Spark Connect".
Se habilitar uma versão de ambiente em um pipeline causar falha, desabilitar essa versão fará com que o pipeline retorne ao seu estado anterior.
Mudanças de comportamento
O Spark Connect possui um pequeno número de diferenças de comportamento em relação ao runtime clássico do PySpark. Consulte Spark Connect vs. classic Spark para a referência completa. O Compatibility scan detecta esses padrões com antecedência e bloqueia a migração até que sejam resolvidos, para que você possa encontrá-los e corrigi-los antes que afetem os dados de produção.
Em um pipeline, as situações mais comuns em que o comportamento pode diferir são:
- Construção DataFrame intercalados e mutação de sessão
- UDFs que fazem referência a um estado mutável do Python
Construção DataFrame intercalados e mutação de sessão
Quando um pipeline constrói um DataFrame e, em seguida, modifica o estado da sessão Spark (por exemplo, altera o catálogo ou esquema default , define uma configuração, substitui uma view temporária ou registra novamente uma UDF), ele usa o DataFrame:
- Sem uma versão de ambiente, o DataFrame utiliza o estado da sessão anterior à mutação .
- Com uma versão de ambiente, o DataFrame utiliza o estado da sessão pós-mutação .
Por exemplo:
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
Sem uma versão de ambiente, mytable contém [(1, "Original Row")]. Com uma versão de ambiente, mytable contém [(2, "Replaced Row")].
UDFs que fazem referência a um estado mutável do Python
Quando uma UDF faz referência a uma variável global do Python cujo valor muda após a definição da UDF:
- Sem uma versão de ambiente, a UDF usa o valor mais recente da variável.
- Com uma versão de ambiente, a UDF usa o valor no momento em que a UDF foi definida .
Por exemplo:
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")))
Sem uma versão de ambiente, my_mv contém [("alex_b",)]. Com uma versão de ambiente, my_mv contém [("alex_a",)].
Se um pipeline depender de algum desses padrões, faça uma auditoria antes de habilitar uma versão do ambiente.
Verificação de compatibilidade
O compatibility scan encontra padrões de código em seu pipeline que produziriam resultados diferentes sob uma versão de ambiente, para que você possa corrigi-los antes que um pipeline seja migrado automaticamente. Quando o scan está habilitado em um pipeline:
- Cada atualização emite um evento
BehaviorChangeInSparkConnectWARNno log de eventos do pipeline por padrão detectado. - O pipeline não é migrado para uma versão de ambiente, e você não pode habilitar uma por conta própria, até que todos os avisos de compatibilidade da atualização bem-sucedida anterior sejam resolvidos.
Esta verificação não se aplica a um pipeline que não tenha atualizações anteriores ou que já tenha uma versão de ambiente definida.
Ative a verificação em um pipeline.
É possível habilitar a verificação de compatibilidade adicionando a configuração de pipeline pipelines.environmentVersion.enableCompatibilityScan. É possível adicionar configurações por meio da interface do usuário do editor de pipeline ou adicionando uma entrada ao JSON de configuração do pipeline.
Através da interface do usuário :
- No editor de pipeline, clique em Configurações .
- Localize a seção Configuração nas configurações do pipeline.
- Clique
Adicionar configuração .
- Insira
pipelines.environmentVersion.enableCompatibilityScancomo key etruecomo valor. - Salve as configurações do pipeline.
No pipeline JSON :
Adicione a seguinte entrada ao bloco configuration :
"configuration": {
"pipelines.environmentVersion.enableCompatibilityScan": "true"
}
Revise e resolva avisos de compatibilidade
Para localizar e limpar os padrões que bloqueiam uma versão de ambiente em seu pipeline:
- Execute o pipeline no modo dry run e, em seguida, faça uma query no log de eventos do pipeline para eventos
BehaviorChangeInSparkConnectWARN. Cada evento relata um padrão detectado. Consulte a referência de eventos de compatibilidade para obter a lista completa de códigos de problemas, exemplos de padrões e correções sugeridas. - Atualize o código do pipeline para remover os padrões detectados seguindo a correção sugerida e execute o pipeline novamente.
- Repita até que uma atualização bem-sucedida não emita mais eventos de compatibilidade. O pipeline pode então ser migrado automaticamente, e você também pode habilitar uma versão de ambiente por conta própria.
A ativação de uma versão de ambiente executa as mesmas verificações de segurança, independentemente de o Databricks migrar o pipeline automaticamente ou de você definir environment_version por conta própria. Um pipeline com avisos de compatibilidade não resolvidos não migra para uma versão de ambiente até que os avisos sejam resolvidos. Se a migração não puder ser concluída com segurança, ou falhar por qualquer motivo, ela será interrompida antes de gravar quaisquer dados e o pipeline continuará a ser executado no runtime anterior.
Quando uma atualização para por um desses motivos, o log de eventos do pipeline e a mensagem de erro da atualização descrevem a causa e os passos para resolvê-la. Siga esses os passos e execute o pipeline novamente para concluir a migração. Se você acredita que um aviso de compatibilidade é um falso positivo, resolva o padrão sinalizado ou entre em contato com o suporte da Databricks.
Referência de eventos de compatibilidade
Quando a verificação de compatibilidade é executada em um pipeline, ela emite um evento BehaviorChangeInSparkConnect WARN no log de eventos do pipeline por padrão detectado. Quando a atualização bem-sucedida anterior detectou quaisquer padrões, o pipeline não é migrado para uma versão do ambiente até que os padrões sejam resolvidos.
Cada evento gera um único código de problema que identifica o que foi detectado. Para consultar um código, localize-o na tabela de códigos de problemas — cada linha contém um link para a seção da categoria que apresenta um padrão de exemplo e a correção sugerida.
Formato do evento
BehaviorChangeInSparkConnect Os eventos seguem o esquema padrão log de eventospipeline:
event_typeébehavior_change_in_spark_connect.leveléWARN.detailscontém o objetobehavior_change_in_spark_connect, que possui um único campoissue. O valor da emissão é um dos códigos listados abaixo.messageÉ uma descrição legível por humanos do padrão detectado.
Códigos de emissão
Categoria | Código do problema | Descrição |
|---|---|---|
| O catálogo default foi alterado após a criação de um DataFrame . O DataFrame existente pode resolver tabelas usando o novo catálogo default . | |
|
| |
| O banco de dados default foi alterado após a criação de um DataFrame . O DataFrame existente pode resolver tabelas usando o novo banco de dados default . | |
|
| |
| A função de fluxo chama um comando de ponto de verificação. | |
| A função de fluxo cria imediatamente uma view DataFrame ( | |
| A função de fluxo cria um perfil de recurso. | |
| A função de fluxo chama | |
| A função de fluxo executa um eager | |
| A função de fluxo realiza uma transação ávida Spark ML . | |
| A função de fluxo registra uma fonte de dados Python . | |
| A função de fluxo opera em um identificador de consulta de transmissão ativo. | |
| A função de fluxo registra ou remove um ouvinte de consulta de transmissão. | |
| A função de fluxo chama | |
| A função de fluxo executa uma | |
| A função de fluxo executa uma | |
| A função de fluxo inicia uma consulta de transmissão ( | |
|
| |
|
| |
|
| |
| Uma view temporária global foi substituída após a criação de um DataFrame que a referenciava. A substituição pode ser refletida no DataFrame existente. | |
| Uma view temporária foi substituída após a criação de um DataFrame que a referenciava. A substituição pode ser refletida no DataFrame existente. | |
| Uma UDF foi registrada novamente com o mesmo nome após a criação de um DataFrame que a referenciava. O DataFrame existente pode usar a nova definição UDF. | |
| Uma UDTF foi registrada novamente com o mesmo nome após a criação de um DataFrame que a referenciava. O DataFrame existente pode usar a nova definição UDTF. | |
| Uma UDF (Função Definida pelo Usuário) faz referência a uma variável global mutável em Python. Com uma versão de ambiente, a UDF usa o valor da variável no momento em que a UDF foi definida, e não no momento da invocação. | |
| Uma UDTF faz referência a uma variável Python global e mutável. Com uma versão de ambiente, a UDTF usa o valor da variável no momento em que a UDTF foi definida, e não no momento da invocação. |
Mutações de banco de dados e catálogo
Esses problemas ocorrem quando o código pipeline modifica o banco de dados ou o catálogo default . Com uma versão de ambiente, os DataFrames construídos antes da mutação podem resolver tabelas usando o novo banco de dados ou catálogo.
Exemplo de padrão que dispara um evento:
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()
Sem uma versão de ambiente, df resolve events do catálogo marketing . Com uma versão de ambiente, df resolve events do catálogo sales .
Sugestão de correção: qualifique totalmente os nomes das tabelas para que a resolução não dependa do catálogo ou banco de dados default e evite alterar o catálogo ou banco de dados default entre a criação e o uso DataFrame .
from pyspark import pipelines as dp
df = spark.read.table("marketing.default.events")
@dp.materialized_view
def events_summary():
return df.groupBy("region").count()
mutações de configuração do Spark
Esses problemas ocorrem quando o código do pipeline modifica a configuração do Spark de maneiras que podem alterar o comportamento do DataFrame em uma determinada versão do ambiente.
Exemplo de padrão que dispara um evento:
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")
Sem uma versão de ambiente, a conversão usa o valor de configuração (conf) no momento da criação do DataFrame. Com uma versão de ambiente, a conversão usa spark.sql.ansi.enabled=true e pode falhar com entrada inválida.
Solução sugerida: Defina todas as configurações necessárias do Spark no início do arquivo de pipeline, antes da criação de qualquer DataFrame. Para configuração por consulta, use a configuração configuration do pipeline na especificação do pipeline.
Substituições temporárias view
Esses problemas ocorrem quando o código pipeline substitui uma view temporária após a criação de um DataFrame que a referencia. Com uma versão de ambiente, o DataFrame existente pode refletir o novo conteúdo view .
Exemplo de padrão que dispara um evento:
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
Sem uma versão de ambiente, mytable contém [(1, "Original Row")]. Com uma versão de ambiente, mytable contém [(2, "Replaced Row")].
Solução sugerida: Crie cada view temporária apenas uma vez e não a substitua. Se você precisar de várias visualizações com dados relacionados, dê a cada uma um nome distinto.
Mutações UDF e UDTF
Esses problemas são gerados quando o código do pipeline modifica uma UDF ou UDTF de maneiras que alteram o comportamento em uma determinada versão do ambiente.
Exemplo de padrão que dispara um evento:
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")))
Sem uma versão de ambiente, my_mv contém [("alex_b",)]. Com uma versão de ambiente, my_mv contém [("alex_a",)].
Solução sugerida: Passe os valores para a UDF como argumentos em vez de capturá-los de variáveis globais do Python, ou defina a variável global antes de definir a UDF e não a modifique posteriormente.
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")))
Execução ágil dentro de funções de fluxo
Esses problemas são emitidos quando o código pipeline executa um comando Spark imediato dentro de uma função decorada por um decorador de pipeline (@table, @materialized_view, etc). Espera-se que as funções de fluxo definam e retornem um DataFrame; Comando ansioso que grava dados, gerencia transmissão de consultas, registro recurso ou execução de operações ML não é permitido dentro de uma função de fluxo com uma versão de ambiente definida.
Sugestão de solução: Mova as operações eager para fora da função de fluxo e retorne um DataFrame da função de fluxo. Efeitos colaterais, como escrever em uma tabela ou iniciar uma consulta de transmissão, estão fora da definição pipeline ; o mecanismo pipeline lida com a materialização do DataFrame retornado pela função de fluxo.
Encontre eventos de compatibilidade no logde eventos.
A consulta a seguir retorna todos os eventos de compatibilidade para um pipeline, ordenados do mais recente para o mais antigo:
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;
Para contabilizar eventos por código de problema em atualizações recentes:
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;
Para saber como consultar o log de eventos, consulte Consultar o logde eventos.
Recursos adicionais
- Configurar versões de ambiente para pipelines — visão geral do recurso, migração automática e como ativar uma versão de ambiente por conta própria.
- Esquema log de eventos do pipeline — esquema completo log eventos pipeline .
- logde eventos do pipeline — como consultar o log de eventos pipeline .