Ingerir arquivos como o tipo FILE
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.
O tipo FILE armazena e query referências a arquivos não estruturados (documentos, imagens e áudio) em tabelas. Esta página mostra como descobrir arquivos, ingeri-los como referências FILE e ingerir novos arquivos incrementalmente à medida que chegam.
Para a referência sobre o tipo FILE, consulte tipoFILE. Para uma visão geral das abordagens para ingestão de dados não estruturados, consulte Tipo de arquivo FILE e dados não estruturados.
FILE as colunas não têm uma ordenação definida. Você não pode usar uma coluna FILE como coluna de partição, coluna de cluster ou key de Z-order.
[[ ## completed ##]] Para obter mais informações, consulte Limites.
Modos de armazenamento
Uma referência FILE pode ser armazenada em um dos dois modos:
FILE EXTERNALfaz referência a arquivos que já existem em um volume do Unity Catalog. O Databricks não oferece suporte ao armazenamento de referênciasFILE EXTERNALpara arquivos armazenados fora de volumes.FILE MANAGEDarmazena cópias de arquivos no armazenamento gerenciado pelo Unity Catalog. Arquivos de fontes fora de volumes, como SharePoint, Google Drive ou SFTP, devem ser ingeridos e armazenados comoFILE MANAGED.
Use list_files para descobrir arquivos
Utilize a list_files table-valued function para descobrir os arquivos disponíveis em um caminho. Ele retorna uma linha por arquivo com seu path, size, modification_time e uma referência FILE:
SELECT * FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');
Para descobrir arquivos em uma origem que requer uma conexão do Unity Catalog, como SharePoint, Google Drive ou SFTP, adicione o parâmetro connection:
SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');
list_files descobre arquivos recursivamente por default. Para saber mais, consulte a função de valor de tabelalist_files.
Ingerir arquivos como referências de ARQUIVO
Selecione uma abordagem de ingestão com base em onde você armazena seus arquivos. Para referenciar arquivos que já estão em um volume do Unity Catalog, use FILE EXTERNAL. Para ingerir arquivos de uma origem externa, copie-os para o armazenamento gerenciado como FILE MANAGED.
Ingira arquivos de volume como FILE EXTERNAL
Para inserir arquivos que já existem em um volume do Unity Catalog, use uma instrução CREATE TABLE AS SELECT (CTAS) com list_files. Isso cria uma tabela com uma coluna FILE EXTERNAL que referencia cada arquivo no local, sem copiar seu conteúdo. O exemplo a seguir cria uma tabela documents com o nome do arquivo, metadados e uma referência FILE para cada arquivo:
CREATE TABLE documents AS
SELECT _metadata.file_name, *
FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');
Ingerir arquivos de origem externos como FILE GERENCIADO
Para gerar referências de FILE para arquivos em uma fonte como SharePoint, Google Drive ou SFTP, ingira os arquivos primeiro e armazene-os como FILE MANAGED. FILE EXTERNAL não é compatível com arquivos armazenados fora de volumes.
O exemplo a seguir ingere arquivos do SharePoint em uma tabela FILE MANAGED:
- SQL
- Python
- Scala
CREATE TABLE managed_documents (
file_name STRING,
path STRING,
size BIGINT,
modification_time TIMESTAMP,
file FILE MANAGED
) USING DELTA
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
INSERT INTO managed_documents
SELECT _metadata.file_name, *
FROM read_files(
'https://example.sharepoint.com/sites/my-site/',
connection => 'my_sharepoint_connection',
format => 'file');
(spark.read.format("file")
.option("databricks.connection", "my_sharepoint_connection")
.load("https://example.sharepoint.com/sites/my-site/")
.selectExpr("_metadata.file_name", "*")
.writeTo("managed_documents").append())
spark.read.format("file")
.option("databricks.connection", "my_sharepoint_connection")
.load("https://example.sharepoint.com/sites/my-site/")
.selectExpr("_metadata.file_name", "*")
.writeTo("managed_documents").append()
Use pipelines para ingerir novos arquivos incrementalmente
Para ingerir novos arquivos conforme eles chegam, use uma tabela de transmissão em um LakeFlow Pipelines que lê a origem com STREAM read_files(..., format => 'file'). Cada atualização de pipeline processa apenas os arquivos adicionados após a última atualização. Veja read_files e Spark Declarative Pipelines.
Para fazer a transmissão incremental de arquivos de uma fonte como o Google Drive:
- Defina o canal do pipeline como
PREVIEW. A ingestão de referênciasFILEem um pipeline requer o canalPREVIEW. - Defina uma tabela de transmissão que leia a fonte com
STREAM read_files(..., format => 'file'), como no código a seguir:
- SQL
- Python
CREATE STREAMING TABLE streaming_documents (
path STRING,
size BIGINT,
modification_time TIMESTAMP,
file FILE MANAGED
)
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/')
AS SELECT *
FROM STREAM read_files(
'https://drive.google.com/drive/folders/my-folder-id',
connection => 'my_gdrive_connection',
format => 'file');
from pyspark import pipelines as dp
@dp.table(
name="streaming_documents",
schema="path STRING, size BIGINT, modification_time TIMESTAMP, file FILE MANAGED",
table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
)
def streaming_documents():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "file")
.option("databricks.connection", "my_gdrive_connection")
.load("https://drive.google.com/drive/folders/my-folder-id")
)
Aplicar atualizações e exclusões com CDC AUTOMÁTICO
Uma ingestão de transmissão adiciona novos arquivos, mas não captura atualizações ou exclusões da origem. Para aplicar essas alterações, leia o feed de alterações da origem com AUTO CDC.
O Databricks recomenda que você primeiro grave os dados de alteração em uma tabela gerenciada, como no exemplo a seguir, e depois aplique AUTO CDC a essa tabela. A aplicação de AUTO CDC diretamente em STREAM read_files(..., readChangeFeed => true) relê o feed de alterações de origem para cada fluxo downstream, o que pode aumentar os custos de processamento.
Ingira o feed de alterações em dois passos. O exemplo a seguir ingere o feed de alterações do SharePoint e, em seguida, aplica-o a uma tabela de transmissão de destino como SCD tipo 1:
- Grave os dados de alteração em uma tabela de transmissão com arquivos gerenciados, como no código a seguir. Defina
readChangeFeed => trueemread_filespara retornar o feed de alteração, que inclui as colunas de metadados_file_id,_sequencee_is_deleted.
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE documents_changes (
_file_id STRING,
_sequence BIGINT,
_is_deleted BOOLEAN,
path STRING,
size BIGINT,
modification_time TIMESTAMP,
file FILE MANAGED
)
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/')
AS SELECT *
FROM STREAM read_files(
'https://example.sharepoint.com/sites/my-site/',
connection => 'my_sharepoint_connection',
format => 'file',
readChangeFeed => true);
from pyspark import pipelines as dp
@dp.table(
name="documents_changes",
table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
)
def documents_changes():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "file")
.option("databricks.connection", "my_sharepoint_connection")
.option("cloudFiles.readChangeFeed", "true")
.load("https://example.sharepoint.com/sites/my-site/")
)
- Use
AUTO CDCpara aplicar as alterações dessa tabela a uma tabela de transmissão de destino, como no código a seguir. Use_file_idcomo a key,_sequencecomo a coluna de sequência e_is_deletedpara identificar exclusões.
- SQL
- Python
CREATE OR REFRESH STREAMING TABLE documents
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
CREATE FLOW documents_cdc AS AUTO CDC INTO
documents
FROM STREAM documents_changes
KEYS (_file_id)
APPLY AS DELETE WHEN _is_deleted = true
SEQUENCE BY _sequence
COLUMNS * EXCEPT (_is_deleted, _sequence)
STORED AS SCD TYPE 1;
from pyspark import pipelines as dp
from pyspark.sql.functions import col, expr
dp.create_streaming_table(
name="documents",
table_properties={"databricks.filespace-preview": "/Volumes/my_catalog/my_schema/filespace/"}
)
dp.create_auto_cdc_flow(
target = "documents",
source = "documents_changes",
keys = ["_file_id"],
sequence_by = col("_sequence"),
apply_as_deletes = expr("_is_deleted = true"),
except_column_list = ["_is_deleted", "_sequence"],
stored_as_scd_type = 1
)
Converter dados binários embutidos em referências de FILE
Se uma tabela já armazenar conteúdos de arquivo como dados binários em linha, use a funçãocreate_file para gravar esses dados no armazenamento e produzir uma referência FILE.
Os exemplos a seguir usam uma tabela gerada pelo usuário, raw_documents, com uma coluna name e uma coluna content que contém os dados binários.
Gravar dados binários em um volume como FILE EXTERNAL
Para gravar os arquivos em um volume do Unity Catalog como arquivos externos, passe um destination_path para create_file, como no código a seguir:
- SQL
- Python
- Scala
CREATE TABLE documents (name STRING, file FILE EXTERNAL) USING DELTA;
INSERT INTO documents (name, file)
SELECT
name,
create_file(
content => content,
destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name
)
FROM raw_documents;
(spark.read.table("raw_documents")
.selectExpr(
"name",
"create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
.writeTo("documents").append())
spark.read.table("raw_documents")
.selectExpr(
"name",
"create_file(content => content, destination_path => '/Volumes/my_catalog/my_schema/my_volume/' || name) AS file")
.writeTo("documents").append()
Gravar dados binários no armazenamento gerenciado como FILE MANAGED
Para armazenar os arquivos como arquivos gerenciados, chame create_file apenas com o conteúdo binário. Quando você omite destination_path, o Unity Catalog faz o upload do conteúdo para o local de armazenamento gerenciado:
- SQL
- Python
- Scala
CREATE TABLE managed_documents (name STRING, file FILE MANAGED) USING DELTA
TBLPROPERTIES ('databricks.filespace-preview' = '/Volumes/my_catalog/my_schema/filespace/');
INSERT INTO managed_documents (name, file)
SELECT name, create_file(content => content)
FROM raw_documents;
(spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append())
spark.read.table("raw_documents")
.selectExpr("name", "create_file(content => content) AS file")
.writeTo("managed_documents").append()
Próximos passos
FILETipo- Tipo de ARQUIVO e dados não estruturados
- Tutorial: criar um pipeline de processamento de arquivos com o tipo FILE
- Saiba mais sobre o Auto Loader. Consulte O que é o Auto Loader?.