ボリューム内の非構造化データを操作する

このページでは、Unity カタログ ボリュームを使用して、非構造化データ ファイルを格納、クエリ、および処理する方法について説明します。 ファイルのアップロード、メタデータのクエリ、AI 関数を使用したファイルの処理、アクセス制御の適用、ボリュームの他の組織との共有を行う方法について説明します。 可能であれば、カタログ エクスプローラー UI を使用してこのチュートリアルを実行する手順が含まれています。 カタログ エクスプローラー オプションが表示されない場合は、指定された Python または SQL コマンドを使用します。

ボリューム機能とユース ケースの完全な概要については、「 Unity カタログ ボリュームとは」を参照してください。

このチュートリアルでは、AI関数を使ってパスごとにファイルを処理します。 ベータ版で利用可能な FILE タイプでは、ファイル参照やメタデータをテーブル内の列値として保存できます。 FILEタイプおよび非構造化データを参照してください。

Requirements

  • Unity カタログが有効になっている Azure Databricks ワークスペース。
  • CREATE CATALOG メタストアに対する特権を持つ。 「カタログの作成」を参照してください。 カタログを作成できない場合は、管理者にアクセスを依頼するか、 CREATE SCHEMA 特権を持つ既存のカタログを使用してください。
  • Databricks Runtime 14.3 LTS 以降。
  • AI 関数の場合: サポートされているリージョン内のワークスペース。
  • OpenSharing の場合: メタストアに対する CREATE SHARE および CREATE RECIPIENT 特権。 「データと AI 資産を安全に共有する」を参照してください。

手順 1: ボリュームを作成する

ファイルを格納するカタログ、スキーマ、ボリュームを作成します。 ボリューム管理の詳細な手順については、「 Unity カタログ ボリュームの作成と管理」を参照してください。

手順 1.1: カタログとスキーマを作成する

SQL

-- Create a catalog
CREATE CATALOG IF NOT EXISTS unstructured_data_lab;
USE CATALOG unstructured_data_lab;

-- Create a schema
CREATE SCHEMA IF NOT EXISTS raw;
USE SCHEMA raw;

Python

spark.sql("CREATE CATALOG IF NOT EXISTS unstructured_data_lab")
spark.sql("USE CATALOG unstructured_data_lab")
spark.sql("CREATE SCHEMA IF NOT EXISTS raw")
spark.sql("USE SCHEMA raw")

カタログ エクスプローラー

  1. [データ] アイコンをクリックします。サイドバーのカタログ
  2. [ 作成>カタログの作成] をクリックします
  3. カタログ名として「unstructured_data_lab」と入力します。
  4. Create をクリックしてください。
  5. [ カタログの表示] をクリックします。

カタログ ページで、次の手順を実行します。

  1. [ スキーマの作成] をクリックします。
  2. rawスキーマ名として入力します。
  3. Create をクリックしてください。

手順 1.2: マネージド ボリュームを作成する

SQL

CREATE VOLUME IF NOT EXISTS files_volume
COMMENT 'Volume for storing unstructured data files';

Python

spark.sql("""
    CREATE VOLUME IF NOT EXISTS files_volume
    COMMENT 'Volume for storing unstructured data files'
""")

カタログ エクスプローラー

スキーマ ページで、次の手順を実行します。

  1. 作成>ボリュームをクリックします。
  2. ボリューム名として「files_volume」と入力します。
  3. [マネージド ボリューム] が選択されていることを確認します。
  4. Create をクリックしてください。

手順 2: ファイルをアップロードする

ファイルをボリュームにアップロードします。 包括的なファイル管理の例については、 Unity カタログ ボリューム内のファイルの操作に関するページを参照してください。

手順 2.1: ファイルをアップロードする

このチュートリアルでは、 databricks-datasets の例を使用することも、カタログ エクスプローラー UI を使用して独自のファイルをアップロードすることもできます。

Python に慣れていない場合でも、Python コマンドを使用して、 databricks-datasets からボリュームにファイルをコピーできます。 ノートブックでコマンドを実行する手順については、「 Databricks ノートブックの管理」を参照してください。

Python

# Upload a single image file
dbutils.fs.cp(
    "dbfs:/databricks-datasets/flower_photos/roses/10090824183_d02c613f10_m.jpg",
    "/Volumes/unstructured_data_lab/raw/files_volume/rose.jpg"
)

# Upload a single PDF file
dbutils.fs.cp(
    "dbfs:/databricks-datasets/COVID/CORD-19/2020-03-13/COVID.DATA.LIC.AGMT.pdf",
    "/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf"
)

# Upload a directory
local_dir = "dbfs:/databricks-datasets/samples/data/mllib"
volume_path = "/Volumes/unstructured_data_lab/raw/files_volume/sample_files"

for file_info in dbutils.fs.ls(local_dir):
    source = file_info.path
    dest = f"{volume_path}/{file_info.name}"
    dbutils.fs.cp(source, dest, recurse=True)
    print(f"Uploaded: {file_info.name}")

カタログ エクスプローラー

[Python] タブの Python コードでは、2 つのファイル (JPG と PDF) と、.txtファイルと.csv ファイルを含むディレクトリがアップロードされます。 カタログ エクスプローラーを使用してファイルをアップロードするには:

  1. ボリューム ページで、[ このボリュームにアップロード] をクリックします。
  2. [ ファイルのアップロード ] ダイアログの [ ファイル] で、[ ファイルの参照 ] をクリックするか、ドロップ ゾーンにファイルをドラッグ アンド ドロップします。
  3. [ 宛先ボリューム] で、前の手順で作成したボリュームが選択されていることを確認します。

手順 2.2: アップロードを確認する

SQL

LIST '/Volumes/unstructured_data_lab/raw/files_volume/';

Python

files = dbutils.fs.ls("/Volumes/unstructured_data_lab/raw/files_volume/")
for f in files:
    print(f"{f.name}\t{f.size} bytes")

カタログ エクスプローラー

ファイルがアップロードされると、ボリューム ページに表示されます。 ファイル名をクリックしてプレビューを表示するか、ディレクトリをクリックして個々のファイルを表示します。

別の方法: %fs マジック コマンドを使用する

%fsマジック コマンドを使用します。

%fs ls /Volumes/unstructured_data_lab/raw/files_volume/

手順 3: ファイル メタデータのクエリを実行する

ファイル情報を照会して、ボリューム内の内容を正確に把握します。 クエリ パターンの詳細については、「SQL を使用した ボリューム内のファイルの一覧表示とクエリ」を参照してください。

手順 3.1: ファイル メタデータを表示する

SQL

SELECT
  path,
  _metadata.file_name,
  _metadata.file_size,
  _metadata.file_modification_time
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile'
);

Python

df = (
    spark.read
    .format("binaryFile")
    .option("recursiveFileLookup", "true")
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")
)

df.select("path", "modificationTime", "length").show(truncate=False)

カタログ エクスプローラー

カタログ エクスプローラーのボリューム ページには、各ファイルの 名前 (拡張子を含む)、 サイズ最終更新日が 表示されます。

手順 4: ファイルのクエリと処理

Azure Databricks AI 関数を使用して、ドキュメントからコンテンツを抽出し、画像を分析します。 AI 関数機能の完全な概要については、AI 関数を使用したデータのエンリッチメントに関するページを参照してください。

AI 関数には、サポートされているリージョン内のワークスペースが必要です。 AI 関数を使用したデータのエンリッチを参照してください。

AI 関数にアクセスできない場合は、代わりに標準の Python ライブラリを使用してください。 例として、以下の代替セクションを展開します。

手順 4.1: ドキュメントを解析する

SQL

SELECT
  path AS file_path,
  ai_parse_document(content, map('version', '2.0')) AS parsed_content
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.pdf'
);

Python

result_df = spark.sql("""
    SELECT
      path AS file_path,
      ai_parse_document(content, map('version', '2.0')) AS parsed_content
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.pdf'
    )
""")
display(result_df)
代替手段: AI 関数を使用せずに PDF を解析する

AI 関数がリージョンで使用できない場合は、Python ライブラリを使用します。

%pip install PyPDF2==3.0.1

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
from PyPDF2 import PdfReader
import io

@udf(returnType=StringType())
def extract_pdf_text(content):
    if content is None:
        return None
    try:
        reader = PdfReader(io.BytesIO(content))
        return "\n".join(page.extract_text() or "" for page in reader.pages)
    except Exception as e:
        return f"Error: {str(e)}"

df = spark.read.format("binaryFile") \
    .option("pathGlobFilter", "*.pdf") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")

result_df = df.withColumn("text_content", extract_pdf_text("content"))
display(result_df.select("path", "text_content"))

手順 4.2: 画像を分析する

SQL

SELECT
  path,
  ai_query(
    'databricks-llama-4-maverick',
    'Describe this image in one sentence:',
    files => content
  ) AS description
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.{jpg,jpeg,png}'
)
WHERE _metadata.file_size < 5000000;

Python

result_df = spark.sql("""
    SELECT
      path,
      ai_query(
        'databricks-llama-4-maverick',
        'Describe this image in one sentence:',
        files => content
      ) AS description
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.{jpg,jpeg,png}'
    )
    WHERE _metadata.file_size < 5000000
""")
display(result_df)
代替手段: AI 関数を使用せずに画像メタデータを抽出する

AI 関数を使用せずに画像メタデータを抽出するには:

%pip install pillow==10.4.0

from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, IntegerType, StringType
from PIL import Image
import io

image_schema = StructType([
    StructField("width", IntegerType()),
    StructField("height", IntegerType()),
    StructField("format", StringType())
])

@udf(returnType=image_schema)
def get_image_info(content):
    if content is None:
        return None
    try:
        img = Image.open(io.BytesIO(content))
        return {"width": img.width, "height": img.height, "format": img.format}
    except:
        return None

df = spark.read.format("binaryFile") \
    .option("pathGlobFilter", "*.{jpg,jpeg,png}") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/")

result_df = df.withColumn("image_info", get_image_info("content"))
display(result_df.select("path", "image_info.*"))

手順 4.3: ファイル名でフィルター処理して分析する

次の使用例は、ファイル名に部分文字列 "rose" を含む画像ファイルをフィルター処理します。

SQL

SELECT
  path AS file_path,
  ai_query(
    'databricks-llama-4-maverick',
    'Describe this image in one sentence:',
    files => content
  ) AS description
FROM read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/',
  format => 'binaryFile',
  fileNamePattern => '*.{jpg,jpeg,png}'
)
WHERE _metadata.file_name ILIKE '%rose%';

Python

result_df = spark.sql("""
    SELECT
      path AS file_path,
      ai_query(
        'databricks-llama-4-maverick',
        'Describe this image in one sentence:',
        files => content
      ) AS description
    FROM read_files(
      '/Volumes/unstructured_data_lab/raw/files_volume/',
      format => 'binaryFile',
      fileNamePattern => '*.{jpg,jpeg,png}'
    )
    WHERE _metadata.file_name ILIKE '%rose%'
""")
display(result_df)

手順 4.4: 構造化テーブルを使用してファイルを結合する

この例では、行番号を使用して、デモンストレーションのためにファイルとタクシー乗車をペアリングします。 運用環境では、意味のあるビジネスキーで結合します。

SQL

-- This example demonstrates joining file metadata with structured data
-- by pairing files with taxi trips using row numbers
WITH files_with_row AS (
  SELECT
    path,
    SPLIT(path, '/')[SIZE(SPLIT(path, '/')) - 1] AS file_name,
    length,
    ROW_NUMBER() OVER (ORDER BY path) AS file_row
  FROM read_files(
    '/Volumes/unstructured_data_lab/raw/files_volume/',
    format => 'binaryFile'
  )
),
trips_with_row AS (
  SELECT
    tpep_pickup_datetime,
    pickup_zip,
    dropoff_zip,
    fare_amount,
    ROW_NUMBER() OVER (ORDER BY tpep_pickup_datetime) AS trip_row
  FROM samples.nyctaxi.trips
  WHERE pickup_zip IS NOT NULL
  LIMIT 5
)
SELECT
  f.path,
  f.file_name,
  f.length,
  t.pickup_zip,
  t.dropoff_zip,
  t.fare_amount,
  t.tpep_pickup_datetime
FROM files_with_row f
INNER JOIN trips_with_row t ON f.file_row = t.trip_row;

Python

from pyspark.sql.functions import col, row_number, element_at, split
from pyspark.sql.window import Window

# Read files and add row numbers
files_df = spark.read.format("binaryFile") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/") \
    .withColumn("file_name", element_at(split(col("path"), "/"), -1))

files_with_row = files_df.alias("files") \
    .withColumn("file_row", row_number().over(Window.orderBy("path")))

# Get trips and add row numbers
trips_df = spark.table("samples.nyctaxi.trips") \
    .filter(col("pickup_zip").isNotNull()) \
    .limit(5)

trips_with_row = trips_df.alias("trips") \
    .withColumn("trip_row", row_number().over(Window.orderBy("tpep_pickup_datetime")))

# Join on row numbers
result_df = files_with_row \
    .join(trips_with_row, col("file_row") == col("trip_row"), "inner") \
    .select(
        "files.path",
        "files.file_name",
        "files.length",
        "trips.pickup_zip",
        "trips.dropoff_zip",
        "trips.fare_amount",
        "trips.tpep_pickup_datetime"
    )

display(result_df)

手順 5: アクセス制御を適用する

ボリューム内のファイルを読み書きできるユーザーを制御します。 Unity カタログでの権限の管理の詳細については、「 Unity カタログでの権限の管理」を参照してください。

手順 5.1: アクセス権を付与する

SQL

-- Replace <user-or-group-name> with your workspace group or user name

-- Grant read access
GRANT READ VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

-- Grant read and write access
GRANT READ VOLUME, WRITE VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

-- Grant all privileges
GRANT ALL PRIVILEGES ON VOLUME unstructured_data_lab.raw.files_volume
TO `<user-or-group-name>`;

Python

# Replace <user-or-group-name> with your workspace group or user name
spark.sql("""
    GRANT READ VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

spark.sql("""
    GRANT READ VOLUME, WRITE VOLUME ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

spark.sql("""
    GRANT ALL PRIVILEGES ON VOLUME unstructured_data_lab.raw.files_volume
    TO `<user-or-group-name>`
""")

カタログ エクスプローラー

  1. ボリューム ページの [ アクセス許可 ] タブに移動します。
  2. [許可] をクリックします。
  3. ユーザーのメール アドレスまたはグループの名前を入力します。
  4. 付与するアクセス許可を選択します。
  5. [Confirm](確認) をクリックします。

手順 5.2: 現在の特権を表示する

SQL

SHOW GRANTS ON VOLUME unstructured_data_lab.raw.files_volume;

Python

display(spark.sql("SHOW GRANTS ON VOLUME unstructured_data_lab.raw.files_volume"))

カタログ エクスプローラー

ボリューム ページの [ アクセス許可 ] タブには、ボリュームにアクセスできるユーザーとグループが表示されます。

手順 6: 増分インジェストを設定する

オートローダーを使用して、ボリュームに到着した新しいファイルを自動的に処理します。 このパターンは、継続的なデータ インジェスト ワークフローに役立ちます。 インジェスト パターンの詳細については、「 一般的なデータ読み込みパターン」を参照してください。

手順 6.1: ストリーミング テーブルを作成する

SQL

CREATE OR REFRESH STREAMING TABLE document_ingestion
SCHEDULE EVERY 1 HOUR
AS SELECT
  path,
  modificationTime,
  length,
  content,
  _metadata,
  current_timestamp() AS ingestion_time
FROM STREAM(read_files(
  '/Volumes/unstructured_data_lab/raw/files_volume/incoming/',
  format => 'binaryFile'
));

Python

from pyspark.sql.functions import current_timestamp, col

dbutils.fs.mkdirs("/Volumes/unstructured_data_lab/raw/files_volume/incoming/")

df = spark.readStream.format("cloudFiles") \
    .option("cloudFiles.format", "binaryFile") \
    .option("pathGlobFilter", "*.pdf") \
    .load("/Volumes/unstructured_data_lab/raw/files_volume/incoming/")

df_enriched = df \
    .withColumn("ingestion_time", current_timestamp()) \
    .withColumn("source_file", col("_metadata.file_path"))

query = df_enriched.writeStream \
    .option("checkpointLocation",
            "/Volumes/unstructured_data_lab/raw/files_volume/_checkpoints/docs") \
    .trigger(availableNow=True) \
    .toTable("document_ingestion")

query.awaitTermination()

手順 7: OpenSharing でファイルを共有する

OpenSharing を使用して、他の組織のユーザーとボリュームを安全に共有します。 共有する前に受信者を作成する必要があります。 受信者は、共有データにアクセスできる外部組織またはユーザーを表します。 受信者のセットアップについては、 OpenSharing のデータ受信者の作成 (Databricks から Databricks への共有) を参照してください。

手順 7.1: 共有を作成して構成する

SQL

-- Create a share
CREATE SHARE IF NOT EXISTS unstructured_data_share
COMMENT 'Document files for partners';

-- Add the volume
ALTER SHARE unstructured_data_share
ADD VOLUME unstructured_data_lab.raw.files_volume;

-- Create a recipient
CREATE RECIPIENT IF NOT EXISTS <partner_org>
USING ID '<recipient-sharing-identifier>';

-- Grant access
GRANT SELECT ON SHARE unstructured_data_share
TO RECIPIENT <partner_org>;

Python

spark.sql("""
    CREATE SHARE IF NOT EXISTS unstructured_data_share
    COMMENT 'Document files for partners'
""")

spark.sql("""
    ALTER SHARE unstructured_data_share
    ADD VOLUME unstructured_data_lab.raw.files_volume
""")

spark.sql("""
    CREATE RECIPIENT IF NOT EXISTS <partner_org>
    USING ID '<recipient-sharing-identifier>'
""")

spark.sql("""
    GRANT SELECT ON SHARE unstructured_data_share
    TO RECIPIENT <partner_org>
""")

手順 7.2: 共有データにアクセスする (受信者として)

SQL

-- View available shares
SHOW SHARES IN PROVIDER <provider_name>;

-- Create a catalog from the share
CREATE CATALOG IF NOT EXISTS shared_documents
FROM SHARE <provider_name>.unstructured_data_share;

-- Query shared files
SELECT * EXCEPT (content), _metadata
FROM read_files(
  '/Volumes/shared_documents/raw/files_volume/',
  format => 'binaryFile'
)
LIMIT 10;

Python

spark.sql("SHOW SHARES IN PROVIDER <provider_name>").show()

spark.sql("""
    CREATE CATALOG IF NOT EXISTS shared_documents
    FROM SHARE <provider_name>.unstructured_data_share
""")

df = spark.read.format("binaryFile") \
    .load("/Volumes/shared_documents/raw/files_volume/")

df.select("path", "modificationTime", "length").show(10)

手順 8: ファイルをクリーンアップする

不要になったらファイルを削除します。

Python

# Delete a single file
dbutils.fs.rm("/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf")

# Delete a directory recursively
dbutils.fs.rm("/Volumes/unstructured_data_lab/raw/files_volume/sample_files/", recurse=True)

CLI

# Delete a single file
databricks fs rm dbfs:/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf

# Delete a directory recursively
databricks fs rm -r dbfs:/Volumes/unstructured_data_lab/raw/files_volume/sample_files/
代替手段: 標準の Python を使用する
import os
os.remove("/Volumes/unstructured_data_lab/raw/files_volume/covid.pdf")

import shutil
shutil.rmtree("/Volumes/unstructured_data_lab/raw/files_volume/sample_files/")

その他のリソース

ボリュームに関する学習を続ける

SQL 関数の参照