Ficheiros de processo com UDFs

Importante

Este recurso está em versão Beta. Os administradores do espaço de trabalho podem controlar o acesso a esse recurso na página Visualizações . Ver Gerir as pré-visualizações de Azure Databricks.

Use uma função definida pelo utilizador (UDF) para processar os ficheiros referenciados por uma FILE coluna com o seu próprio código e bibliotecas. A UDF recebe cada FILE valor como referência de ficheiro nativa da língua. Pode ler os bytes do ficheiro ou abri-lo como um caminho local, depois devolver um valor de metadados, um ficheiro derivado ou uma saída transformada.

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

Ler metadados de ficheiros num UDF

Um FILE valor tem campos de metadados que podes ler sem abrir o ficheiro. A tabela seguinte contém os campos disponíveis:

Acessor Description
uri O URI do ficheiro.
offset Um deslocamento para o ficheiro, em bytes.
size O tamanho do arquivo, em bytes.
content_type O tipo MIME do ficheiro, quando conhecido.
checksum Uma soma de verificação usada para identificar a versão do ficheiro, como <algorithm>:<value>.

Aceda a estes campos com notação pontual no FILE valor, como mostrado no código seguinte:

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 ficheiro numa UDF

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

  • as_local_file(): Devolve um caminho local que pode passar a qualquer biblioteca que aceite um caminho de ficheiro, como uma biblioteca de imagens ou multimédia.
  • open(): Devolve um fluxo binário que lê apenas os bytes que solicita, em vez de materializar o ficheiro completo.

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

Extrair as dimensões da imagem

Pode usar um UDF escalar para devolver as dimensões de uma imagem como uma width x height cadeia. O UDF pede as_local_file() para obter um caminho local, depois passa esse caminho para uma biblioteca de imagens padrão (PILem Python, ImageIO no Scala), como mostrado no seguinte código:

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

Detetar o tipo de um ficheiro a partir dos seus bytes

O UDF seguinte lê apenas os primeiros oito bytes de cada ficheiro com open() e deteta o tipo de ficheiro a partir do seu número mágico, sem materializar o ficheiro completo:

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 ficheiros com uma tabela UDF (UDTF)

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

Quando uma UDTF (ou qualquer UDF) escreve novos ficheiros com FileRef.from_bytes, o seu código deve cumprir os seguintes requisitos:

  • Cria o volume alvo antes de executares o UDTF. Um trabalhador Python não consegue criar um volume de topo de nível. 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.
  • Passa por um caminho absoluto dbfs: . Devolver a FileRef 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. Apague ou ignore ficheiros que já existam antes de escrever. Como FileRef.from_bytes escreve com flags exclusive-create, escrever sobre um ficheiro existente levanta FileExistsError.

Exemplo: Extrair fotogramas de vídeo

O seguinte UDTF lê um vídeo FILE, extrai cada fotograma com a av biblioteca (PyAV), grava-o num volume e gera uma linha por fotograma:

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 numa linha por frame:

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 GOVERN FILE com filtros de linha

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

Filtro de linha

Um filtro de linha é um UDF que devolve um BOOLEAN. As linhas para as quais devolve false são omitidas dos resultados das consultas.

O filtro de linhas seguinte mantém apenas linhas com ficheiros que referenciam uma folha de cálculo Excel, com base nos metadados do content_type ficheiro:

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

Registe um UDF no Catálogo Unity

Registar um UDF de processamento de ficheiros no Catálogo Unity para o governar com permissões de catálogo e reutilizá-lo em cadernos, consultas e utilizadores. Registar e executar um UDF requer os seguintes privilégios:

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

O exemplo seguinte regista um UDF SQL que devolve a extensão de um ficheiro, 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 registar um UDF Python ou Scala no Catálogo Unity, consulte as funções definidas pelo utilizador (UDFs) SQL e Python no Catálogo Unity e as funções de tabela definidas pelo utilizador (UDTFs) em Python no Catálogo Unity.

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

O código UDF corre com os privilégios do proprietário da função, não do chamador da função. Os privilégios do proprietário aplicam-se à leitura dos bytes de um FILE. Um chamador com apenas EXECUTE permissões no UDF, e sem acesso direto ao volume subjacente, pode ainda assim desencadear leituras dos ficheiros referenciados.

Como um UDF de processamento de ficheiros é um caminho de acesso regulado ao conteúdo dos ficheiros, considere os seguintes efeitos secundários de segurança e governação:

  • Os utilizadores podem aceder ao conteúdo dos ficheiros através do UDF. Conceda EXECUTE permissões apenas a utilizadores que pretenda dar acesso indireto ao conteúdo dos ficheiros.
  • Quem chama herda o acesso ao ficheiro do proprietário. Verifique se o proprietário da UDF tem acesso a volumes não mais amplo do que o que os chamadores deveriam ter.

Para mais informações sobre como o Azure Databricks determina o utilizador autorizado à medida que a execução passa para um corpo UDF, veja Utilizador autorizado e utilizador de sessão.

Passos seguintes