Processar arquivos com UDFs
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 |
|---|---|
| O URI do arquivo. |
| Um offset no arquivo, em bytes. |
| O tamanho do arquivo, em bytes. |
| O tipo MIME do arquivo, quando conhecido. |
| Uma soma de verificação usada para identificar a versão do arquivo, como |
Acesse estes campos com notação de ponto no valor FILE, conforme mostrado no código a seguir:
- Python
- Scala
- SQL
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()
import org.apache.spark.sql.functions.{col, udf}
val isLargeImage = udf { (file: FileRef) =>
file.contentType.startsWith("image/") && file.size > 5000000L
}
spark.read.table("documents").select(col("file.uri"), isLargeImage(col("file"))).display()
SELECT file.uri, file.content_type, file.size
FROM documents
WHERE file.content_type LIKE 'image/%'
AND file.size > 5000000;
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
- Scala
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()
import org.apache.spark.sql.functions.{col, udf}
import javax.imageio.ImageIO
val imageResolution = udf { (file: FileRef) =>
// asLocalFile() returns a java.io.File.
val image = ImageIO.read(file.asLocalFile())
s"${image.getWidth}x${image.getHeight}"
}
spark.read.table("images").select(col("photo.uri"), imageResolution(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
- Scala
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()
import org.apache.spark.sql.functions.{col, udf}
val fileSignature = udf { (file: FileRef) =>
// open() returns a java.io.InputStream.
val stream = file.open()
try {
val header = new Array[Byte](8)
val n = stream.read(header)
if (n >= 4 && header(0) == '%' && header(1) == 'P' && header(2) == 'D' && header(3) == 'F') "pdf"
else if (n >= 4 && header(0) == 0x89.toByte && header(1) == 'P' && header(2) == 'N' && header(3) == 'G') "png"
else if (n >= 3 && header(0) == 0xFF.toByte && header(1) == 0xD8.toByte && header(2) == 0xFF.toByte) "jpeg"
else "unknown"
} finally {
stream.close()
}
}
spark.read.table("documents").select(col("file.uri"), fileSignature(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 umFileRefpara uma tabela Delta Lake requer um URIdbfs:, comodbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Um caminho simples geraDELTA_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_bytesgrava com sinalizadores de criação exclusiva, gravar sobre um arquivo existente geraFileExistsError.
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:
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:
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
- Python
- Scala
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);
from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType
@udf(returnType=BooleanType())
def excel_only(file):
return file.content_type in (
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"application/vnd.ms-excel")
import org.apache.spark.sql.functions.udf
val excelOnly = udf { (file: FileRef) =>
Set(
"application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
"application/vnd.ms-excel").contains(file.contentType)
}
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:
USAGEeCREATEno esquema, eUSAGEno catálogo. - Para execução de uma UDF:
EXECUTEna UDF eUSAGEno 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:
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
EXECUTEapenas 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
FILETipoai_parse_documentfunção- Tipo de Arquivo
- Funções definidas pelo usuário (UDFs) em SQL e Python no Unity Catalog
- Usuário autorizado e usuário da sessão
- Aplique manualmente filtros de linha e máscaras de coluna
- Funções escalares definidas pelo usuário (UDFs) em Python
- Funções de tabela definidas pelo usuário (UDTFs) em Python
- Tipo de ARQUIVO e dados não estruturados