Dosyaları FILE tipi olarak alın

Important

Bu özellik Beta sürümündedir. Çalışma alanı yöneticileri Bu özelliğe erişimi Önizlemeler sayfasından denetleyebilir. Bkz. Azure Databricks önizlemelerini yönetme.

Tür, FILE yapılandırılmamış dosyalara (belgeler, görseller ve ses) referansları tablolarda depolar ve sorgular. Bu sayfa, dosyaların nasıl keşfedileceğini, referans FILE olarak nasıl alınacağını ve yeni dosyalara ulaştıkça kademeli olarak nasıl alınacağını gösterir.

Tiple FILE ilgili referans için bkz.FILE Yapılandırılmamış veri alımına yönelik yaklaşımların genel bir özeti için bkz. FILE tipi ve yapılandırılmamış veri.

Uyarı

FILE Sütunların belirli bir sıralaması yoktur. Bir FILE sütunu bölümleme sütunu, küme sütunu veya Z-order anahtarı olarak kullanamazsınız. Daha fazla bilgi için bkz . Sınırlar.

Depolama modları

Bir FILE referans iki moddan birinde saklanabilir:

  • FILE EXTERNAL Unity Kataloğu ciltinde zaten var olan referans dosyaları. Databricks, hacimler dışında depolanan dosyalar için referansların saklanmasını FILE EXTERNAL desteklemez.
  • FILE MANAGED dosya kopyalarını Unity Kataloğu tarafından yönetilen depolamada depolar. SharePoint, Google Drive veya SFTP gibi hacim dışındaki kaynaklardan gelen dosyalar, alınmalı ve saklanmalıdırFILE MANAGED.

Dosyaları keşfetmek için kullanın list_files

list_files Tablo değerli fonksiyonu kullanarak bir yolda mevcut dosyaları keşfet. Her dosya için bir satır döner, , pathsizemodification_time, ve bir FILE referans içerir:

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

SharePoint, Google Drive veya SFTP gibi Unity Kataloğu bağlantısı gerektiren bir kaynakta dosyaları keşfetmek için şu parametreyi connection ekleyin:

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

list_files Dosyaları varsayılan olarak özyinelemeli olarak keşfeder. Daha fazla bilgi için tablo değerli fonksiyona bakınızlist_files.

Dosyaları FILE referansı olarak alın

Dosyalarınızı nerede sakladığınıza göre bir alım yöntemi seçin. Zaten bir Unity Kataloğu ciltinde FILE EXTERNALbulunan dosyalara referans vermek için . Dosyaları harici bir kaynaktan almak için onları yönetilen depolamaya kopyalayın.FILE MANAGED

Hacim dosyalarını DOSYA HARİCİ olarak alın

Bir Unity Kataloğu ciltinde zaten var olan dosyaları almak için (CTAS) ifadesi ile list_filesbir (CTAS) ifadesi CREATE TABLE AS SELECT kullanın. Bu, FILE EXTERNAL her dosyayı yerinde referans veren bir sütunlu bir tablo oluşturur ve içeriğini kopyalamaz. Aşağıdaki örnek, dosya adı, meta veri ve her dosya için bir FILE referans içeren bir documents tablo oluşturur:

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

Dış kaynak dosyaları DOSYA YÖNETİLİLİRİ olarak alın

SharePoint, Google Drive veya SFTP gibi bir kaynaktaki dosyalar için referans oluşturmak FILE için önce dosyaları alın ve onları . FILE MANAGEDşeklinde depolayın. FILE EXTERNAL hacimler dışında depolanan dosyalar için desteklenmiyor.

Aşağıdaki örnek, SharePoint'ten dosyaları bir FILE MANAGED tabloya alır:

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

Yeni dosyaları kademeli olarak almak için boru hatlarını kullanın

Yeni dosyalar geldikçe alın mı, bir Lakeflow boru hattında kaynağı okuyan bir akış tablosu kullanın.STREAM read_files(..., format => 'file') Her boru hattı güncellemesi yalnızca son güncellemeden sonra eklenen dosyaları işliyor. Bkz read_files . ve Spark Declarative Pipeline.

Google Drive gibi bir kaynaktan dosyaları kademeli olarak akış yapmak için:

  1. Boru hattının kanalını .PREVIEW Bir boru hattında referans almak FILE için kanal PREVIEW gereklidir.

  2. Aşağıdaki kodda olduğu gibi, kaynağı okuyan STREAM read_files(..., format => 'file')bir akış tablosu tanımlayın:

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

AUTO CDC ile güncellemeler ve silme işlemlerini uygulayın

Bir akış girişi yeni dosyalar ekler ama kaynaktan güncellemeleri veya silmeleri yakalamaz. Bu değişiklikleri uygulamak için, kaynak değişiklik akışını .AUTO CDC

Warning

Databricks, değişiklik verilerini önce yönetilen bir tabloya yerleştirmenizi önerir, aşağıdaki örnekte olduğu gibi, ardından o tabloya uygularsınız AUTO CDC . AUTO CDC Doğrudan uygulamak, STREAM read_files(..., readChangeFeed => true) her aşağı akış için kaynak değişim beslemesini yeniden okunur, bu da işlem maliyetlerini artırabilir.

Değişim akışını iki adımda alın. Aşağıdaki örnek, değişim akışını SharePoint'ten alır ve ardından SCD tip 1 olarak hedef akış tablosuna uygular:

  1. Değişiklik verilerini, aşağıdaki kodda olduğu gibi yönetilen dosyalarla bir akış tablosuna yazın. Değişim akışını döndürmek için ayarlanın readChangeFeed => trueread_files; bu besleme , _file_id_sequence, ve _is_deleted meta veri sütunlarını içerir.

    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. Aşağıdaki kodda olduğu gibi, o tablodan gelen değişiklikleri hedef akış tablosuna uygulamak için kullanılır AUTO CDC . Anahtar olarak, _sequence sıralı sütun olarak ve _is_deleted silinmeleri tanımlamak için kullanın_file_id.

    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
    )
    

Hat içi ikili veriyi FILE referanslarına dönüştürün

Bir tablo dosya içeriğini zaten satır içi ikili veri olarak saklıyorsa, create_file fonksiyon kullanarak bu veriyi depolamaya yazıp bir FILE referans üretin.

Aşağıdaki örnekler, bir sütun ve content ikili veriyi barındıran bir sütun içeren kullanıcı tarafından oluşturulan bir tablo raw_documentskullanırname.

Bir hacme DOSYA HARİCİ olarak ikili veri yazmak

Dosyaları bir Unity Kataloğu hacmine harici dosya olarak yazmak için, aşağıdaki kodda olduğu gibi a'yı destination_path 'ye create_filegeçirin:

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

İkili verileri yönetilen depolamaya DOSYA YÖNETİLİLİK olarak yazmak

Dosyaları yönetilen dosya olarak saklamak için sadece ikili içerikle çağrı create_file yapın. Atladığınızda destination_path, Unity Catalog içeriği yönetilen depolama konumuna yükler:

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

Sonraki Adımlar