Memasukkan file sebagai tipe FILE

Important

Fitur ini ada di Beta. Admin ruang kerja dapat mengontrol akses ke fitur ini dari halaman Pratinjau . Lihat Kelola Pratinjau Azure Databricks.

Tipe menyimpan FILE dan mengurai referensi ke file tidak terstruktur (dokumen, gambar, dan audio) dalam tabel. Halaman ini menunjukkan cara menemukan file, mengingkatnya sebagai FILE referensi, dan secara bertahap mengingini file baru saat tiba.

Untuk referensi pada FILE tipe, lihat FILE tipe. Untuk gambaran umum pendekatan dalam mengingsipasi data tidak terstruktur, lihat FILE type dan data unstruktur.

Catatan

FILE kolom tidak memiliki urutan yang jelas. Anda tidak dapat menggunakan FILE kolom sebagai kolom partisi, kolom klaster, atau kunci urutan Z. Untuk informasi selengkapnya, lihat Batasan.

Mode penyimpanan

Sebuah FILE referensi dapat disimpan dalam salah satu dari dua mode:

  • FILE EXTERNAL merujuk file yang sudah ada di volume Unity Catalog. Databricks tidak mendukung penyimpanan FILE EXTERNAL referensi untuk file yang disimpan di luar volume.
  • FILE MANAGED menyimpan salinan file di penyimpanan yang dikelola oleh Unity Catalog. File dari sumber di luar volume, seperti SharePoint, Google Drive, atau SFTP, harus dimasukkan dan disimpan sebagai FILE MANAGED.

Gunakan list_files untuk menemukan file

Gunakan list_files fungsi bernilai tabel untuk menemukan file yang tersedia di sebuah jalur. Ia mengembalikan satu baris per file dengan path, size, modification_time, dan referensi:FILE

SELECT * FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Untuk menemukan file di sumber yang memerlukan koneksi Unity Catalog, seperti SharePoint, Google Drive, atau SFTP, tambahkan parameter berikutconnection:

SELECT * FROM list_files('https://example.sharepoint.com/sites/my-site/', connection => 'my_sharepoint_connection');

list_files secara default menemukan file secara rekursif. Untuk mempelajari lebih lanjut, lihat list_files fungsi bernilai tabel.

Memasukkan file sebagai referensi FILE

Pilih pendekatan penghapusan berdasarkan tempat Anda menyimpan file Anda. Untuk merujuk file yang sudah ada di volume Unity Catalog, gunakan FILE EXTERNAL. Untuk mengingest file dari sumber eksternal, salin ke penyimpanan terkelola sebagai FILE MANAGED.

Masukkan file volume sebagai FILE EKSTERNAL

Untuk mengingest file yang sudah ada dalam volume Unity Catalog, gunakan CREATE TABLE AS SELECT pernyataan (CTAS) dengan list_files. Ini membuat tabel dengan FILE EXTERNAL kolom yang merujuk setiap file di tempatnya, tanpa menyalin isinya. Contoh berikut membuat documents tabel dengan nama file, metadata, dan FILE referensi untuk setiap file:

CREATE TABLE documents AS
  SELECT _metadata.file_name, *
  FROM list_files('/Volumes/my_catalog/my_schema/raw_files/');

Masukkan file sumber eksternal sebagai FILE MANAGED

Untuk menghasilkan FILE referensi file di sumber seperti SharePoint, Google Drive, atau SFTP, masukkan file terlebih dahulu dan simpan sebagai FILE MANAGED. FILE EXTERNAL tidak didukung untuk file yang disimpan di luar volume.

Contoh berikut menginput file dari SharePoint ke dalam sebuah FILE MANAGED tabel:

SQL

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');

Python

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

Scala

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

Gunakan pipeline untuk mengingest file baru secara bertahap

Untuk menginget file baru saat tiba, gunakan tabel streaming dalam pipeline Lakeflow yang membaca sumber dengan STREAM read_files(..., format => 'file'). Setiap pembaruan pipeline hanya memproses file yang ditambahkan setelah pembaruan terakhir. Lihat read_files dan Spark Pipeline Deklaratif.

Untuk melakukan streaming file secara bertahap dari sumber seperti Google Drive:

  1. Atur saluran pipeline ke PREVIEW. Menginesting FILE referensi dalam pipeline memerlukan saluran tersebut PREVIEW .

  2. Definisikan tabel streaming yang membaca sumber dengan STREAM read_files(..., format => 'file'), seperti dalam kode berikut:

    SQL

    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');
    

    Python

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

Terapkan pembaruan dan penghapusan dengan AUTO CDC

Streaming ingest menambahkan file baru tetapi tidak menangkap pembaruan atau penghapusan dari sumbernya. Untuk menerapkan perubahan tersebut, baca feed perubahan sumber dengan AUTO CDC.

Warning

Databricks merekomendasikan agar Anda menempatkan data perubahan terlebih dahulu di tabel terkelola, seperti pada contoh berikut, lalu terapkan AUTO CDC ke tabel tersebut. Menerapkan AUTO CDC langsung untuk STREAM read_files(..., readChangeFeed => true) membaca ulang sumber perubahan feed untuk setiap aliran hilir, yang mungkin meningkatkan biaya pemrosesan.

Masukkan feed perubahan dalam dua langkah. Contoh berikut menginput feed perubahan dari SharePoint, lalu menerapkannya ke tabel streaming target sebagai tipe SCD 1:

  1. Tulis data perubahan ke dalam tabel streaming dengan file yang dikelola, seperti pada kode berikut. Atur readChangeFeed => true untuk read_files mengembalikan feed perubahan, yang mencakup _file_idkolom , _sequence, dan _is_deleted metadata.

    SQL

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

    Python

    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/")
      )
    
  2. Gunakan AUTO CDC untuk menerapkan perubahan dari tabel tersebut ke tabel streaming target, seperti dalam kode berikut. Gunakan _file_id sebagai kunci, _sequence sebagai kolom urutan, dan _is_deleted untuk mengidentifikasi penghapusan.

    SQL

    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;
    

    Python

    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
    )
    

Konversi data biner inline menjadi referensi FILE

Jika sebuah tabel sudah menyimpan isi file sebagai data biner inline, gunakan create_file fungsi untuk menulis data tersebut ke penyimpanan dan menghasilkan referensi.FILE

Contoh berikut menggunakan tabel yang dihasilkan pengguna, raw_documents, dengan kolom name dan kolom content yang memuat data biner.

Tulis data biner ke volume sebagai FILE EXTERNAL

Untuk menulis file ke volume Unity Catalog sebagai file eksternal, berikan a destination_path ke create_file, seperti dalam kode berikut:

SQL

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;

Python

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

Scala

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

Tulis data biner ke penyimpanan terkelola sebagai FILE MANAGED

Untuk menyimpan file sebagai file yang dikelola, panggil create_file hanya dengan konten biner. Ketika Anda menghilangkan destination_path, Unity Catalog mengunggah konten ke lokasi penyimpanan terkelola:

SQL

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;

Python

(spark.read.table("raw_documents")
  .selectExpr("name", "create_file(content => content) AS file")
  .writeTo("managed_documents").append())

Scala

spark.read.table("raw_documents")
  .selectExpr("name", "create_file(content => content) AS file")
  .writeTo("managed_documents").append()

Langkah berikutnya