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 EXTERNALUnity カタログのボリュームに既に存在するファイルの参照ファイルです。 Databricksは、ボリューム外に保存されたファイルのFILE EXTERNAL参照の保存をサポートしていません。
Azure Databricksファイルレベルの権限と組み込みのコンプライアンスを活用するワークロードにはFILE MANAGEDを推奨しています。 ガバナンスとライフサイクルの挙動の比較については、 FILEタイプと非構造化データを参照してください。
list_filesを使ってファイルを発見しましょう
list_files table-value関数 table-value関数を使って、パス上の利用可能なファイルを発見します。 ファイルごとに1行を返し、 path、 size、 modification_time、 FILE 参照を示します。
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ドライブのようなソースからファイルを段階的にストリーミングするには:
パイプラインのチャンネルを
PREVIEWに設定してください。 パイプライン内でFILE参照を取り込むにはPREVIEWチャネルが必要です。以下のコードのように、
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として適用します。
変更データを管理ファイル付きのストリーミングテーブルに書き込みます。以下のコードの通りです。
readChangeFeed => trueread_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/") )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_path を create_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()
次のステップ
-
FILE型 - ファイルタイプと非構造化データ
- チュートリアル:ファイル形式でファイル処理パイプラインを構築する
- 自動ローダーの詳細については、こちらを参照してください。 「自動ローダーとは」を参照してください。