Arquivos de processo com UDFs

Importante

Esse recurso está em Beta. Os administradores do workspace podem controlar o acesso a esse recurso na página Visualizações . Consulte Gerenciar visualizações do Azure Databricks.

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

Esta página mostra UDFs de processamento de arquivos em Python, Scala e SQL. Para a FILE referência tipográfica, veja FILE tipo. Para autoria geral de UDF, veja funções definidas pelo usuário (UDFs) escalares em Python, UDFs Scala e Java com escopo de sessão, e funções de tabela definidas pelo usuário em Python (UDTFs).

Leia metadados de arquivo em um UDF

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

Acessor Description
uri O URI do arquivo.
offset Um deslocamento para dentro do arquivo, em bytes.
size O tamanho do arquivo, em bytes.
content_type O tipo MIME do arquivo, quando conhecido.
checksum Um checksum usado para identificar a versão do arquivo, como <algorithm>:<value>.

Acesse esses campos com notação pontual no FILE valor, 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()

Scala

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()

SQL

SELECT file.uri, file.content_type, file.size
  FROM documents
  WHERE file.content_type LIKE 'image/%'
    AND file.size > 5000000;

Leia o conteúdo do arquivo em uma UDF

Um FILE valor possui 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 mídia.
  • open(): Retorna um fluxo binário que lê apenas os bytes que você solicita, em vez de materializar o arquivo inteiro.

Ambos exigem computação do Azure Databricks (um notebook ou um worker UDF) e não estão disponíveis em um cliente Azure Databricks Connect. Você pode declarar FILE como um parâmetro UDF ou tipo de retorno em Python, Scala e SQLs UDFs. Para a API completa, veja FileType.

Extrair as dimensões da imagem

Você pode usar um UDF escalar para retornar as dimensões de uma imagem como uma width x height string. O UDF chama as_local_file() para obter um caminho local, então passa esse caminho para uma biblioteca de imagens padrão (PILem Python, ImageIO em Scala), como 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()

Scala

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

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()

Scala

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 múltiplos arquivos com um UDF de tabela (UDTF)

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

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

  • Crie o volume alvo antes de rodar o UDTF. Um trabalhador em 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() podem criar subdiretórios, mas não o volume em si.
  • Passe por um caminho absoluto dbfs: . Retornar a FileRef a a a uma tabela Delta Lake requer um dbfs: URI, como dbfs:/Volumes/my_catalog/my_schema/frames/frame_00000.jpg. Um caminho nu eleva DELTA_VIOLATE_CONSTRAINT_WITH_VALUES.
  • Verificar que as escritas são idempotentes. Exclua ou pule arquivos que já existem antes de escrever. Como FileRef.from_bytes escreve com flags de criação exclusiva, escrever sobre um arquivo existente gera FileExistsError.

Exemplo: extrair quadros de vídeo

O seguinte UDTF lê um vídeo FILE, extrai cada quadro com a av biblioteca (PyAV), grava em um volume e gera 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 alvo com uma FILE EXTERNAL coluna, depois chame o 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;

Colunas de FILE de governar com filtros de linha

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

Filtro de linhas

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

O filtro de linhas a seguir mantém apenas linhas com arquivos que referenciam uma planilha Excel, com base nos metadados do content_type 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);

Python

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")

Scala

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 como aplicar e gerenciar filtros de linhas, incluindo os passos e limitações do Explorador de Catálogos, veja Aplicar manualmente filtros de linha e máscaras de coluna.

Registre um UDF no Catálogo Unity

Registre um UDF de processamento de arquivos no Catálogo Unity para governá-lo com permissões de catálogo e reutilize-o em cadernos, consultas e usuários. Registrar e rodar um UDF requer os seguintes privilégios:

  • Para criar um UDF: USAGE e CREATE no esquema, e USAGE no catálogo.
  • Para rodar um UDF: EXECUTE no UDF, e USAGE no esquema e catálogo.

O exemplo a seguir registra um UDF SQL que retorna a extensão de um arquivo, depois chama o 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 um UDF em Python ou Scala no Catálogo Unity, veja funções definidas pelo usuário (UDFs) em SQL e Python no Catálogo Unity e funções de tabela definidas pelo usuário (UDTFs) em Python no Catálogo Unity.

Segurança: UDFs funcionam com os privilégios do proprietário

O código UDF roda 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 EXECUTE permissões no UDF, e sem acesso direto ao volume subjacente, ainda pode acionar leituras dos arquivos referenciados.

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

  • Os usuários podem acessar o conteúdo dos arquivos usando o UDF. Conceda EXECUTE permissões apenas a usuários que você pretende dar acesso indireto ao conteúdo dos arquivos.
  • Quem liga herda o acesso ao arquivo do proprietário. Verifique se o proprietário da UDF tem acesso a volumes não mais amplo do que os chamantes deveriam ter.

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

Próximas Etapas