Pular para o conteúdo principal

Compatibilidade de versão do ambiente

info

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 ...") e createOrReplaceTempView.
  • Usa APIs PySpark que não estão disponíveis no Spark Connect, incluindo SparkContext, RDD, SQLContext e 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

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:

Python
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:

Python
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 BehaviorChangeInSparkConnect WARN no 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 :

  1. No editor de pipeline, clique em Configurações .
  2. Localize a seção Configuração nas configurações do pipeline.
  3. Clique Ícone de mais (+). Adicionar configuração .
  4. Insira pipelines.environmentVersion.enableCompatibilityScan como key e true como valor.
  5. Salve as configurações do pipeline.

No pipeline JSON :

Adicione a seguinte entrada ao bloco configuration :

JSON
"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:

  1. Execute o pipeline no modo dry run e, em seguida, faça uma query no log de eventos do pipeline para eventos BehaviorChangeInSparkConnect WARN. 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.
  2. Atualize o código do pipeline para remover os padrões detectados seguindo a correção sugerida e execute o pipeline novamente.
  3. 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.
  • details contém o objeto behavior_change_in_spark_connect , que possui um único campo issue . 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

Mutações de banco de dados e catálogo

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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 .

Mutações de banco de dados e catálogo

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE CATALOG foi chamada para fora de um evento decorado por um decorador de tubulações. O catálogo default pode mudar inesperadamente em operações subsequentes.

Mutações de banco de dados e catálogo

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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 .

Mutações de banco de dados e catálogo

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE DATABASE foi chamada para fora de um evento decorado por um decorador de tubulações. O banco de dados default pode mudar inesperadamente em operações subsequentes.

Execução ágil dentro de funções de fluxo

CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama um comando de ponto de verificação.

Execução ágil dentro de funções de fluxo

CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo cria imediatamente uma view DataFrame (createOrReplaceTempView ou similar).

Execução ágil dentro de funções de fluxo

CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo cria um perfil de recurso.

Execução ágil dentro de funções de fluxo

GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama spark.resources ou uma API de recurso relacionada.

Execução ágil dentro de funções de fluxo

MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa um eager MERGE INTO em uma tabela de destino.

Execução ágil dentro de funções de fluxo

ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo realiza uma transação ávida Spark ML .

Execução ágil dentro de funções de fluxo

REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo registra uma fonte de dados Python .

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo opera em um identificador de consulta de transmissão ativo.

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo registra ou remove um ouvinte de consulta de transmissão.

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama spark.streams para gerenciar consultas de transmissão.

Execução ágil dentro de funções de fluxo

WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa uma DataFrameWriterV2 operações ansiosas.

Execução ágil dentro de funções de fluxo

WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa uma DataFrame.write operações ansiosas.

Execução ágil dentro de funções de fluxo

WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo inicia uma consulta de transmissão (writeStream.start()).

mutações de configuração do Spark

CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED

spark.conf.set() ou spark.conf.unset() foi chamado dentro de uma função decorada por um decorador de pipeline. Isso não é compatível com uma versão de ambiente.

mutações de configuração do Spark

SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.set() foi chamada fora de uma função decorada com um decorador de pipeline após a criação de um DataFrame . A alteração de configuração pode afetar o DataFrame existente em tempo de execução.

mutações de configuração do Spark

UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.unset() foi chamada fora de uma função decorada com um decorador de pipeline após a criação de um DataFrame . A alteração de configuração pode afetar o DataFrame existente em tempo de execução.

Substituições temporárias view

REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Substituições temporárias view

REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Categoria

Código do problema

Descrição

Mutações de banco de dados e catálogo

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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 .

Mutações de banco de dados e catálogo

USE_CATALOG_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE CATALOG foi chamada para fora de um evento decorado por um decorador de tubulações. O catálogo default pode mudar inesperadamente em operações subsequentes.

Mutações de banco de dados e catálogo

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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 .

Mutações de banco de dados e catálogo

USE_DATABASE_OUTSIDE_QUERY_FUNCTION_COULD_CHANGE_BEHAVIOR

USE DATABASE foi chamada para fora de um evento decorado por um decorador de tubulações. O banco de dados default pode mudar inesperadamente em operações subsequentes.

Execução ágil dentro de funções de fluxo

CHECKPOINT_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama um comando de ponto de verificação.

Execução ágil dentro de funções de fluxo

CREATE_DATAFRAME_VIEW_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo cria imediatamente uma view DataFrame (createOrReplaceTempView ou similar).

Execução ágil dentro de funções de fluxo

CREATE_RESOURCE_PROFILE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo cria um perfil de recurso.

Execução ágil dentro de funções de fluxo

GET_RESOURCES_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama spark.resources ou uma API de recurso relacionada.

Execução ágil dentro de funções de fluxo

MERGE_INTO_TABLE_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa um eager MERGE INTO em uma tabela de destino.

Execução ágil dentro de funções de fluxo

ML_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo realiza uma transação ávida Spark ML .

Execução ágil dentro de funções de fluxo

REGISTER_DATA_SOURCE_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo registra uma fonte de dados Python .

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo opera em um identificador de consulta de transmissão ativo.

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_LISTENER_BUS_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo registra ou remove um ouvinte de consulta de transmissão.

Execução ágil dentro de funções de fluxo

STREAMING_QUERY_MANAGER_COMMAND_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo chama spark.streams para gerenciar consultas de transmissão.

Execução ágil dentro de funções de fluxo

WRITE_OPERATION_V2_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa uma DataFrameWriterV2 operações ansiosas.

Execução ágil dentro de funções de fluxo

WRITE_OPERATION_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo executa uma DataFrame.write operações ansiosas.

Execução ágil dentro de funções de fluxo

WRITE_STREAM_OPERATION_START_WITHIN_QUERY_FUNCTION_NOT_SUPPORTED

A função de fluxo inicia uma consulta de transmissão (writeStream.start()).

mutações de configuração do Spark

CHANGE_CONF_INSIDE_QUERY_FUNCTION_NOT_SUPPORTED

spark.conf.set() ou spark.conf.unset() foi chamado dentro de uma função decorada por um decorador de pipeline. Isso não é compatível com uma versão de ambiente.

mutações de configuração do Spark

SET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.set() foi chamada fora de uma função decorada com um decorador de pipeline após a criação de um DataFrame . A alteração de configuração pode afetar o DataFrame existente em tempo de execução.

mutações de configuração do Spark

UNSET_CONF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

spark.conf.unset() foi chamada fora de uma função decorada com um decorador de pipeline após a criação de um DataFrame . A alteração de configuração pode afetar o DataFrame existente em tempo de execução.

Substituições temporárias view

REPLACE_GLOBAL_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Substituições temporárias view

REPLACE_TEMP_VIEW_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

OVERWRITE_SESSION_UDF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

OVERWRITE_SESSION_UDTF_AFTER_DATAFRAME_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

UDF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

Mutações UDF e UDTF

UDTF_REFERENCES_GLOBAL_VARIABLE_COULD_CHANGE_BEHAVIOR

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.

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:

Python
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 .

Python
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:

Python
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:

Python
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:

Python
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.

Python
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:

SQL
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:

SQL
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