Pular para o conteúdo principal

Processar arquivos com UDFs

info

Beta

Este recurso está em Beta. Os administradores do Workspace podem controlar o acesso a este recurso na página Pré-visualizações . Consulte Gerenciar prévias do Databricks.

Use uma função definida pelo usuário (UDF) para processar os arquivos referenciados por uma coluna FILE com seu próprio código e bibliotecas. A UDF recebe cada valor FILE como uma referência de arquivo nativa da linguagem. Ele pode ler os bytes do arquivo ou abri-lo como um caminho local e, em seguida, retornar um valor de metadados, um arquivo derivado ou uma saída transformada.

Esta página mostra UDFs de processamento de arquivos em Python, Scala e SQL. Para a referência do tipo FILE, consulte tipoFILE. Para a criação geral de UDFs, consulte Funções escalares definidas pelo usuário (UDFs) em Python, UDFs Scala e Java com escopo de sessão e Funções de tabela definidas pelo usuário (UDTFs) em Python.

Ler metadados de arquivo em um UDF

Um valor FILE possui campos de metadados que você pode ler sem abrir o arquivo. A tabela a seguir contém os campos disponíveis:

Acessor

Descrição

uri

O URI do arquivo.

offset

Um offset no arquivo, em bytes.

size

O tamanho do arquivo, em bytes.

content_type

O tipo MIME do arquivo, quando conhecido.

checksum

Uma soma de verificação usada para identificar a versão do arquivo, como <algorithm>:<value>.

Acessor

Descrição

uri

O URI do arquivo.

offset

Um offset no arquivo, em bytes.

size

O tamanho do arquivo, em bytes.

content_type

O tipo MIME do arquivo, quando conhecido.

checksum

Uma soma de verificação usada para identificar a versão do arquivo, como <algorithm>:<value>.

Acesse estes campos com notação de ponto no valor FILE, conforme mostrado no código a seguir:

Python
from pyspark.sql.functions import col, udf
from pyspark.sql.types import BooleanType

@udf(returnType=BooleanType())
def is_large_image(file):
return file.content_type.startswith("image/") and file.size > 5_000_000

spark.read.table("documents").select(col("file").uri, is_large_image(col("file"))).display()

Ler o conteúdo do arquivo em um UDF

Um valor FILE tem dois métodos para ler o arquivo subjacente:

  • as_local_file(): Retorna um caminho local que você pode passar para qualquer biblioteca que aceite um caminho de arquivo, como uma biblioteca de imagens ou de mídia.
  • open(): Retorna uma transmissão binária que lê apenas os bytes solicitados, em vez de materializar o arquivo inteiro.

Ambos exigem compute do Databricks (um notebook ou um worker de UDF) e não estão disponíveis em um cliente Databricks Connect. Você pode declarar FILE como um parâmetro de UDF ou tipo de retorno em UDFs Python, Scala e SQL. Para a API completa, consulte Tipo de Arquivo.

Extrair dimensões da imagem

Você pode usar uma UDF escalar para retornar as dimensões de uma imagem como uma string width x height. A UDF chama as_local_file() para obter um caminho local e, em seguida, passa esse caminho para uma biblioteca de imagens padrão (PIL em Python, ImageIO em Scala), conforme mostrado no código a seguir:

Python
from pyspark.sql.functions import col, udf
from pyspark.sql.types import StringType
from PIL import Image

@udf(returnType=StringType())
def image_resolution(file):
# as_local_file() returns a pathlib.Path.
with Image.open(file.as_local_file()) as img:
return f"{img.width}x{img.height}"

spark.read.table("images").select(col("photo").uri, image_resolution(col("photo"))).display()

Detectar o tipo de um arquivo a partir de seus bytes

O UDF a seguir lê apenas os primeiros oito bytes de cada arquivo com open() e detecta o tipo de arquivo a partir de seu número mágico, sem materializar o arquivo inteiro:

Python
from pyspark.sql.functions import col, udf
from pyspark.sql.types import StringType

@udf(returnType=StringType())
def file_signature(file):
with file.open() as f:
header = f.read(8)
if header.startswith(b"%PDF"):
return "pdf"
if header.startswith(b"\x89PNG"):
return "png"
if header.startswith(b"\xff\xd8\xff"):
return "jpeg"
return "unknown"

spark.read.table("documents").select(col("file").uri, file_signature(col("file"))).display()

Gerar vários arquivos com uma UDF de tabela (UDTF)

Para transformar um arquivo de entrada em vários arquivos de saída, como ao dividir um vídeo em quadros, use uma UDF de tabela (UDTF). A UDTF recebe um FILE como entrada e gera uma linha por arquivo de saída, criando cada arquivo com FileRef.from_bytes(). Declare a coluna de arquivo como FILE no esquema returnType da UDTF. Para a criação geral de UDTFs, consulte Funções de tabela definidas pelo usuário (UDTFs) em Python.

Quando uma UDTF (ou qualquer UDF) grava novos arquivos com FileRef.from_bytes, seu código deve atender aos seguintes requisitos:

  • Crie o volume de destino antes da execução do UDTF. Um worker Python não pode criar um volume de nível superior. Crie-o com CREATE VOLUME IF NOT EXISTS. Dentro de um volume existente, os.makedirs() pode criar subdiretórios, mas não o volume em si.
  • Passe um caminho dbfs: absoluto. Retornar um FileRef para uma tabela Delta Lake requer um URI dbfs:, como dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Um caminho simples gera DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Verifique se as gravações são idempotentes. Exclua ou ignore arquivos que já existem antes de gravar. Como FileRef.from_bytes grava com sinalizadores de criação exclusiva, gravar sobre um arquivo existente gera FileExistsError.

Exemplo: Extrair quadros de vídeo

A UDTF a seguir lê um vídeo FILE, extrai cada quadro com a biblioteca av (PyAV), grava-o em um volume e retorna uma linha por quadro:

Python
import io
import os
import av
from pyspark.sql.functions import udtf
from pyspark.sql.types import FileRef

@udtf(returnType="clip_id STRING, frame_index INT, frame FILE")
class ExtractFrames:
def __init__(self):
self.output_dir = "/Volumes/my_catalog/my_schema/frames/"
os.makedirs(self.output_dir, exist_ok=True)

def eval(self, video):
clip_id = video.uri.split("/")[-1].split(".")[0]
container = av.open(video.as_local_file())
stream = container.streams.video[0]
for i, frame in enumerate(container.decode(stream)):
buffer = io.BytesIO()
frame.to_image().save(buffer, format="JPEG")

local_path = os.path.join(self.output_dir, f"{clip_id}_frame_{i:05d}.jpg")
if os.path.exists(local_path):
os.remove(local_path)

yield (
clip_id,
i,
FileRef.from_bytes(buffer.getvalue(), path=f"dbfs:{local_path}", content_type="image/jpeg"),
)
container.close()

spark.udtf.register("extract_frames", ExtractFrames)

Crie a tabela de destino com uma coluna FILE EXTERNAL, em seguida, chame a UDTF com LATERAL para expandir cada vídeo em uma linha por quadro:

SQL
CREATE TABLE my_catalog.my_schema.drive_frames (
clip_id STRING,
frame_index INT,
frame FILE EXTERNAL
);

INSERT INTO my_catalog.my_schema.drive_frames
SELECT *
FROM my_catalog.my_schema.drive_clips AS c
JOIN LATERAL extract_frames(c.video) AS f;

Controlar colunas FILE com filtros de linha

Governe uma coluna FILE com filtros de linha com base na identidade do chamador ou nos metadados do arquivo.

Filtro de linha

Um filtro de linha é uma UDF que retorna um BOOLEAN. As linhas para as quais retorna false são omitidas dos resultados da query.

O filtro de linha a seguir mantém apenas as linhas com arquivos que fazem referência a uma planilha do Excel, com base nos metadados content_type do arquivo:

SQL
CREATE FUNCTION excel_only(file FILE)
RETURN file.content_type IN (
'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet',
'application/vnd.ms-excel');

ALTER TABLE documents SET ROW FILTER excel_only ON (file);

Para mais informações sobre a aplicação e o gerenciamento de filtros de linha, incluindo os passos do Explorador de Catálogos e as limitações, consulte Aplicar manualmente filtros de linha e máscaras de coluna.

Registro de uma UDF no Unity Catalog

Registre um UDF de processamento de arquivos no Unity Catalog para governá-lo com permissões de catálogo e reutilizá-lo em notebooks, queries e usuários. O registro e a execução de um UDF exigem os seguintes privilégios:

  • Para criar uma UDF: USAGE e CREATE no esquema, e USAGE no catálogo.
  • Para execução de uma UDF: EXECUTE na UDF e USAGE no esquema e no catálogo.

O exemplo a seguir registra uma UDF SQL que retorna a extensão de um arquivo e, em seguida, chama a UDF para criar uma nova coluna:

SQL
CREATE FUNCTION my_catalog.my_schema.file_extension(file FILE)
RETURNS STRING
RETURN lower(element_at(split(file.uri, '\\.'), -1));

SELECT file.uri, my_catalog.my_schema.file_extension(file) AS extension
FROM documents;

Para registrar uma UDF Python ou Scala no Unity Catalog, consulte Funções definidas pelo usuário (UDFs) SQL e Python no Unity Catalog e Funções de tabela definidas pelo usuário (UDTFs) Python no Unity Catalog.

Segurança: UDFs têm execução com os privilégios do proprietário

A execução do código da UDF ocorre com os privilégios do proprietário da função, não do chamador da função. Os privilégios do proprietário se aplicam à leitura dos bytes de um FILE. Um chamador com apenas permissões EXECUTE na UDF, e sem acesso direto ao volume subjacente, ainda pode Trigger leituras dos arquivos referenciados.

Como um UDF de processamento de arquivos é um caminho de acesso governado ao conteúdo de arquivos, considere os seguintes efeitos colaterais de segurança e governança:

  • Os usuários podem acessar o conteúdo de arquivos usando o UDF. Conceda permissões EXECUTE apenas aos usuários aos quais você pretende dar acesso indireto ao conteúdo de arquivos.
  • Os chamadores herdam o acesso a arquivos do proprietário. Verifique se o proprietário do UDF tem acesso ao volume não mais amplo do que o que os chamadores devem ter.

Para obter mais informações sobre como o Databricks determina o usuário autorizado à medida que a execução cruza para um corpo de UDF, consulte Usuário autorizado e usuário da sessão.

Próximos passos