ファイルをファイル形式として取り込む

Important

この機能は ベータ版です。 ワークスペース管理者は、[ プレビュー] ページからこの機能へのアクセスを制御できます。 Manage Azure Databricks プレビューを参照してください。

FILEタイプは、非構造化ファイル(文書、画像、音声)への参照をテーブルに保存・クエリします。 このページでは、ファイルの発見方法、 FILE 参照として取り込み、新しいファイルが届くたびに段階的に取り込む方法を示しています。

FILE型に関する参考文献については、FILEを参照してください。 非構造化データの取り込み方法の概要については、 FILE type および unstructured data を参照してください。

Note

FILE 列には明確な順序はありません。 FILE列を分割列、クラスタリング列、Z順序キーとして使うことはできません。 詳細については、「制限事項」をご覧ください。

ストレージ モード

FILE参照は2つのモードのいずれかで格納できます。

  • FILE MANAGEDファイルのコピーをUnity Catalog管理ストレージに保存します:権限はテーブルを通じて管理され、行を削除すると参照ファイルがガベージコレクション対象となり、テーブルとそのファイルが同期されたままになります。SharePoint、Googleドライブ、SFTPなど、ボリューム以外のソースからのファイルはFILE MANAGEDとして取り込み保存する必要があります。
  • FILE EXTERNAL Unity カタログのボリュームに既に存在するファイルの参照ファイルです。 Databricksは、ボリューム外に保存されたファイルの FILE EXTERNAL 参照の保存をサポートしていません。

Azure Databricksファイルレベルの権限と組み込みのコンプライアンスを活用するワークロードにはFILE MANAGEDを推奨しています。 ガバナンスとライフサイクルの挙動の比較については、 FILEタイプと非構造化データを参照してください。

list_filesを使ってファイルを発見しましょう

list_files table-value関数 table-value関数を使って、パス上の利用可能なファイルを発見します。 ファイルごとに1行を返し、 pathsizemodification_timeFILE 参照を示します。

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

SharePoint、Google Drive、SFTPなどUnityカタログ接続が必要なソース内のファイルを検出するには、connectionパラメータを追加してください:

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

list_files デフォルトで再帰的にファイルを発見します。 詳細は list_files テーブル値関数をご覧ください。

ファイルをファイル参照として取り込む

ファイルを保存する場所に基づいて取り込み方法を選択してください。 外部ソースからファイルを取り込むには、 FILE MANAGEDとして管理ストレージにコピーしてください。 Unityカタログのボリュームにあるファイルをコピーせずに参照するには FILE EXTERNALを使います。

外部ソースファイルをファイル管理として取り込む

SharePoint、Googleドライブ、SFTPなどのソース内のファイルのFILE参照を生成するには、まずファイルを取り込み、FILE MANAGEDとして保存してください。 FILE EXTERNAL ボリューム外に保存されたファイルには対応していません。

以下の例は、SharePointのファイルをFILE MANAGEDテーブルに取り込みます:

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

ボリュームファイルをファイル外部として取り込む

Unityカタログボリュームに既に存在するファイルを取り込むには、CREATE TABLE AS SELECT付きのlist_files(CTAS)文を使用します。 これにより、各ファイルの内容をコピーせずに各ファイルを参照する、 FILE EXTERNAL 列を持つテーブルが作成されます。 以下の例は、ファイル名、メタデータ、各ファイルのdocuments参照を含むFILEテーブルを作成します。

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

パイプラインを使って新しいファイルを段階的に取り込みます

新しいファイルが届いたときに取り込むには、Lakeflowパイプライン内のストリーミングテーブルを使い、ソースを STREAM read_files(..., format => 'file')で読み取ってください。 各パイプライン更新は、前回の更新後に追加されたファイルのみを処理します。 read_filesおよびSpark宣言的パイプラインを参照してください。

Googleドライブのようなソースからファイルを段階的にストリーミングするには:

  1. パイプラインのチャンネルを PREVIEWに設定してください。 パイプライン内で FILE 参照を取り込むには PREVIEW チャネルが必要です。

  2. 以下のコードのように、 STREAM read_files(..., format => 'file')でソースを読み取るストリーミングテーブルを定義します。

    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で更新や削除を適用してください

ストリーミングインジェストは新しいファイルを追加しますが、ソースからの更新や削除は記録しません。 これらの変更を適用するには、ソースの変更フィードを読み込んで AUTO CDCしてください。

Warnung

Databricksは、まず変更データを管理されたテーブルにランディングし、そのテーブルに適用 AUTO CDC することを推奨しています。 AUTO CDCを直接STREAM read_files(..., readChangeFeed => true)に適用すると、各下流フローのソース変更フィードを再読み込み、処理コストが増加する可能性があります。

変更フィードは2ステップで取り込みます。 以下の例は、SharePointから変更フィードを取り込み、それをターゲットストリーミングテーブルにSCDタイプ1として適用します。

  1. 変更データを管理ファイル付きのストリーミングテーブルに書き込みます。以下のコードの通りです。 readChangeFeed => true read_filesを設定して、_file_id_sequence_is_deletedメタデータの列を含む変更フィードを返します。

    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. AUTO CDCを使って、そのテーブルから変更をターゲットストリームテーブルに適用してください。以下のコードのように。 _file_idをキー、_sequenceをシーケンス列に、_is_deletedで欠失を特定します。

    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
    )
    

インラインバイナリデータをFILE参照に変換する

テーブルがすでにファイルの内容をインラインバイナリデータとして保存している場合は、関数create_fileそのデータをストレージに書き込み、FILE参照を生成します。

以下の例は、raw_documents列とname列を持つユーザー生成のテーブルcontentを使用しています。

管理ストレージにバイナリデータを「ファイル管理」として書き込む

管理ファイルとしてファイルを保存するには、バイナリコンテンツのみで create_file を呼び出してください。 destination_pathを省略すると、Unity Catalogが管理されたストレージにコンテンツをアップロードします:

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

ファイル外部としてボリュームにバイナリデータを書き込む

Unity カタログのボリュームに外部ファイルとして書き込むには、以下のコードのように destination_pathcreate_fileに渡します。

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

次のステップ