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. O UDF recebe cada valor FILE como uma referência de arquivo nativa da linguagem. Em Python, essa referência é um objeto FileRef que você pode importar de pyspark.sql.types. O UDF 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, FileRef

@udf(returnType=BooleanType())
def is_large_image(file: FileRef) -> bool:
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.

nota

Uma UDF registrada no Unity Catalog (CREATE FUNCTION) pode ler os metadados de um FILE, mas não o conteúdo, e não pode criar arquivos. Use uma UDF com escopo de sessão para ler o conteúdo de um arquivo (open, as_local_file) ou criar um (from_bytes, from_local_file).

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 FileRef, StringType
from PIL import Image

@udf(returnType=StringType())
def image_resolution(file: FileRef) -> str:
# 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 FileRef, StringType

@udf(returnType=StringType())
def file_signature(file: FileRef) -> str:
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: FileRef):
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.

Próximos passos​