メインコンテンツまでスキップ

Data Import Hands-On

これは、データ準備や collection のセットアップから実際のデータインポート処理まで、Zilliz Cloud でのデータインポートをすばやく開始できるようにする短期集中コースです。このチュートリアル全体を通して、以下を学びます。

  • スキーマを定義し、対象の collection をセットアップする方法

  • BulkWriter を使用してソースデータを準備し、リモートストレージ bucket に書き込む方法

  • bulk-import API を呼び出してデータをインポートする方法

📘注意

Zilliz Cloud では現在、cluster をホストしているクラウドプロバイダーに関係なく、任意のオブジェクトストレージサービスから任意の Zilliz Cloud cluster にデータをインポートできます。たとえば、AWS S3 bucket から GCP 上にデプロイされた Zilliz Cloud cluster にデータをインポートできます。

低レイテンシかつ安定した体験を確保するために、対象 cluster と同じプロバイダーで、同じリージョンにある bucket または blob container を使用することを推奨します。

始める前に

スムーズに進めるために、以下のセットアップを完了していることを確認してください。

Zilliz Cloud cluster をセットアップする

  • まだの場合は、cluster を作成してください。

  • 次の情報を控えておいてください: Cluster EndpointAPI KeyCluster ID

依存関係をインストールする

現在、データインポート関連の API は Python または Java で使用できます。

Python API を使用するには、ターミナルで次のコマンドを実行して pymilvusminio をインストールするか、最新バージョンにアップグレードしてください。

shell
python3 -m pip install --upgrade pymilvus minio

リモートストレージ bucket を構成する

  • AWS S3 を使用してリモート bucket をセットアップします。

  • 次の情報を控えておいてください。

    • S3 互換ブロックストレージサービスの場合: Access KeySecret KeyBucket Name

    • Microsoft Azure blob ストレージサービスの場合: AccountNameAccountKeyContainerName

    これらの情報は、bucket をホストしているクラウドプロバイダーのコンソールで確認できます。

サンプルコードをより使いやすくするために、構成情報を保存するための変数を使用することを推奨します。

python
## The value of the URL is fixed.
CLOUD_API_ENDPOINT = "https://api.cloud.zilliz.com"
API_KEY=""

# Configs for Zilliz Cloud cluster
CLUSTER_ENDPOINT=""
CLUSTER_ID="" # Zilliz Cloud cluster ID, like "in01-xxxxxxxxxxxxxxx"
COLLECTION_NAME="zero_to_hero"

# Configs for remote bucket
BUCKET_NAME=""
ACCESS_KEY=""
SECRET_KEY=""

対象 collection のスキーマをセットアップする

上記の出力に基づいて、対象 collection のスキーマを作成できます。

次のデモでは、事前定義されたスキーマには最初の 4 つのフィールドを含め、残りの 4 つは動的フィールドとして使用します。

python
from pymilvus import MilvusClient, DataType

# You need to work out a collection schema out of your dataset.
schema = MilvusClient.create_schema(
auto_id=False,
enable_dynamic_field=True
)

DIM = 512

schema.add_field(field_name="id", datatype=DataType.INT64, is_primary=True),
schema.add_field(field_name="bool", datatype=DataType.BOOL),
schema.add_field(field_name="int8", datatype=DataType.INT8),
schema.add_field(field_name="int16", datatype=DataType.INT16),
schema.add_field(field_name="int32", datatype=DataType.INT32),
schema.add_field(field_name="int64", datatype=DataType.INT64),
schema.add_field(field_name="float", datatype=DataType.FLOAT),
schema.add_field(field_name="double", datatype=DataType.DOUBLE),
schema.add_field(field_name="varchar", datatype=DataType.VARCHAR, max_length=512),
schema.add_field(field_name="json", datatype=DataType.JSON),
schema.add_field(field_name="array_str", datatype=DataType.ARRAY, max_capacity=100, element_type=DataType.VARCHAR, max_length=128)
schema.add_field(field_name="array_int", datatype=DataType.ARRAY, max_capacity=100, element_type=DataType.INT64)
schema.add_field(field_name="float_vector", datatype=DataType.FLOAT_VECTOR, dim=DIM),
schema.add_field(field_name="binary_vector", datatype=DataType.BINARY_VECTOR, dim=DIM),
schema.add_field(field_name="float16_vector", datatype=DataType.FLOAT16_VECTOR, dim=DIM),
# schema.add_field(field_name="bfloat16_vector", datatype=DataType.BFLOAT16_VECTOR, dim=DIM),
schema.add_field(field_name="sparse_vector", datatype=DataType.SPARSE_FLOAT_VECTOR)

schema.verify()

print(schema)

上記コード内のパラメータは以下のとおりです。

  • fields:

    • id は primary field です。

    • float_vector は浮動小数点 vector field です。

    • binary_vector は binary vector field です。

    • float16_vector は半精度浮動小数点 vector field です。

    • sparse_vector は疎 vector field です。

    • 残りのフィールドは scalar field です。

  • auto_id=False

    これはデフォルト値です。これを True に設定すると、BulkWriter は生成されるファイルに primary field を含めなくなります。

  • enable_dynamic_field=True

    この値のデフォルトは False です。これを True に設定すると、BulkWriter は生成されるファイルに未定義フィールドとその値をキーと値のペアとして含め、それらを $meta という予約済み JSON フィールドに配置できます。

スキーマの設定が完了したら、次のように対象 collection を作成できます。

python
from pymilvus import MilvusClient

# 1. Set up a Milvus client
client = MilvusClient(
uri=CLUSTER_ENDPOINT,
token=API_KEY
)

# 2. Set index parameters
index_params = MilvusClient.prepare_index_params()

index_params.add_index(
field_name="float_vector",
index_type="AUTOINDEX",
metric_type="IP"
)

index_params.add_index(
field_name="binary_vector",
index_type="AUTOINDEX",
metric_type="HAMMING"
)

index_params.add_index(
field_name="float16_vector",
index_type="AUTOINDEX",
metric_type="IP"
)

index_params.add_index(
field_name="sparse_vector",
index_type="AUTOINDEX",
metric_type="IP"
)

# 3. Create collection
client.create_collection(
collection_name=COLLECTION_NAME,
schema=schema,
index_params=index_params
)

ソースデータを準備する

BulkWriter は、データセットを JSON、Parquet、または NumPy ファイルに書き換えることができます。ここでは RemoteBulkWriter を作成し、この writer を使ってデータをこれらの形式に書き換えます。

RemoteBulkWriter を作成する

schema の準備ができたら、その schema を使って RemoteBulkWriter を作成できます。RemoteBulkWriter は、リモート bucket へのアクセス権限を必要とします。リモート bucket にアクセスするための接続パラメータを ConnectParam オブジェクトに設定し、それを RemoteBulkWriter で参照する必要があります。

python
from pymilvus.bulk_writer import RemoteBulkWriter, BulkFileType
# Use `from pymilvus import RemoteBulkWriter, BulkFileType`
# if your pymilvus version is earlier than 2.4.2

# Connections parameters to access the remote bucket
conn = RemoteBulkWriter.S3ConnectParam(
endpoint="s3.amazonaws.com", # Use "storage.googleapis.com" for Google Cloud Storage
access_key=ACCESS_KEY,
secret_key=SECRET_KEY,
bucket_name=BUCKET_NAME,
secure=True
)
📘メモ

endpoint パラメータは、クラウドプロバイダのストレージサービス URI を指します。

S3 互換ストレージサービスの場合、使用可能な URI は次のとおりです。

  • s3.amazonaws.com(AWS S3)

  • storage.googleapis.com (GCS)

Azure blob storage container の場合は、以下のような有効な接続文字列を使用する必要があります。

DefaultEndpointsProtocol=https;AccountName=<accountName>;AccountKey=<accountKey>;EndpointSuffix=core.windows.net

次に、以下のように RemoteBulkWriter で接続パラメータを参照できます。

python
writer = RemoteBulkWriter(
schema=schema, # Target collection schema
remote_path="/", # Output directory relative to the remote bucket root
segment_size=1024*1024*1024, # Maximum segment size when segmenting the raw data
connect_param=conn, # Connection parameters defined above
file_type=BulkFileType.PARQUET # Type of the generated file.
)

# Possible file types:
# - BulkFileType.JSON,
# - BulkFileType.NPY, and
# - BulkFileType.PARQUET

上記の writer は、JSON 形式でファイルを生成し、指定した bucket のルートフォルダにアップロードします。

  • remote_path="/"

    これは、リモート bucket 内で生成されたファイルの出力パスを決定します。

    "/" に設定すると、RemoteBulkWriter は生成されたファイルをリモート bucket のルートフォルダに配置します。別のパスを使用するには、リモート bucket のルートからの相対パスを設定してください。

  • file_type=BulkFileType.PARQUET

    これは生成されるファイルの種類を決定します。指定可能な値は次のとおりです。

    • BulkFileType.JSON

    • BulkFileType.PARQUET

    • BulkFileType.NPY

  • segment_size=1024*1024*1024

    これは、BulkWriter が生成されたファイルを分割するかどうかを決定します。デフォルト値は 1024 MB (1024 * 1024 * 1024) です。データセットに大量のレコードが含まれる場合は、segment_size を適切な値に設定してデータを分割することをおすすめします。

writer を使用する

writer には 2 つのメソッドがあります。1 つはソースデータセットから行を追加するためのもので、もう 1 つはデータをリモートファイルにコミットするためのものです。

以下のように、ソースデータセットから行を追加できます。

python
import random, string, json
import numpy as np
import tensorflow as tf

def generate_random_str(length=5):
letters = string.ascii_uppercase
digits = string.digits

return ''.join(random.choices(letters + digits, k=length))

# optional input for binary vector:
# 1. list of int such as [1, 0, 1, 1, 0, 0, 1, 0]
# 2. numpy array of uint8
def gen_binary_vector(to_numpy_arr):
raw_vector = [random.randint(0, 1) for i in range(DIM)]
if to_numpy_arr:
return np.packbits(raw_vector, axis=-1)
return raw_vector

# optional input for float vector:
# 1. list of float such as [0.56, 1.859, 6.55, 9.45]
# 2. numpy array of float32
def gen_float_vector(to_numpy_arr):
raw_vector = [random.random() for _ in range(DIM)]
if to_numpy_arr:
return np.array(raw_vector, dtype="float32")
return raw_vector

# # optional input for bfloat16 vector:
# # 1. list of float such as [0.56, 1.859, 6.55, 9.45]
# # 2. numpy array of bfloat16
# def gen_bf16_vector(to_numpy_arr):
# raw_vector = [random.random() for _ in range(DIM)]
# if to_numpy_arr:
# return tf.cast(raw_vector, dtype=tf.bfloat16).numpy()
# return raw_vector

# optional input for float16 vector:
# 1. list of float such as [0.56, 1.859, 6.55, 9.45]
# 2. numpy array of float16
def gen_fp16_vector(to_numpy_arr):
raw_vector = [random.random() for _ in range(DIM)]
if to_numpy_arr:
return np.array(raw_vector, dtype=np.float16)
return raw_vector

# optional input for sparse vector:
# only accepts dict like {2: 13.23, 45: 0.54} or {"indices": [1, 2], "values": [0.1, 0.2]}
# note: no need to sort the keys
def gen_sparse_vector(pair_dict: bool):
raw_vector = {}
dim = random.randint(2, 20)
if pair_dict:
raw_vector["indices"] = [i for i in range(dim)]
raw_vector["values"] = [random.random() for _ in range(dim)]
else:
for i in range(dim):
raw_vector[i] = random.random()
return raw_vector

for i in range(2000):
writer.append_row({
"id": np.int64(i),
"bool": True if i % 3 == 0 else False,
"int8": np.int8(i%128),
"int16": np.int16(i%1000),
"int32": np.int32(i%100000),
"int64": np.int64(i),
"float": np.float32(i/3),
"double": np.float64(i/7),
"varchar": f"varchar_{i}",
"json": json.dumps({"dummy": i, "ok": f"name_{i}"}),
"array_str": np.array([f"str_{k}" for k in range(5)], np.dtype("str")),
"array_int": np.array([k for k in range(10)], np.dtype("int64")),
"float_vector": gen_float_vector(True),
"binary_vector": gen_binary_vector(True),
"float16_vector": gen_fp16_vector(True),
# "bfloat16_vector": gen_bf16_vector(True),
"sparse_vector": gen_sparse_vector(True),
f"dynamic_{i}": i,
})
if (i+1)%1000 == 0:
writer.commit()
print('committed')

print(writer.batch_files)

writer の append_row() メソッドは、行ディクショナリを受け取ります。

行ディクショナリには、schema で定義されたすべてのフィールドをキーとして含める必要があります。dynamic fields が許可されている場合は、未定義のフィールドを含めることもできます。詳細は、Use BulkWriter を参照してください。

BulkWriter は、commit() メソッドを呼び出した後にのみファイルを生成します。

python
writer.commit()

ここまでで、BulkWriter は指定したリモート bucket 内にソースデータを準備しました。

生成されたファイルを確認するには、writer の data_path プロパティを出力して実際の出力パスを取得できます。

python
print(writer.data_path)

# /5868ba87-743e-4d9e-8fa6-e07b39229425
📘メモ

BulkWriter は UUID を生成し、指定された出力ディレクトリ内にその UUID を使ったサブフォルダを作成し、生成されたすべてのファイルをそのサブフォルダ内に配置します。

詳細は、Use BulkWriter を参照してください。

準備したデータをインポートする

このステップの前に、準備済みデータが目的の bucket にすでにアップロードされていることを確認してください。

インポートを開始する

準備したソースデータをインポートするには、以下のように bulk_import() 関数を呼び出す必要があります。

python
from pymilvus.bulk_writer import bulk_import

# Publicly accessible URL for the prepared data in the remote bucket
object_url = "s3://{0}/{1}/".format(BUCKET_NAME, str(writer.data_path)[1:])
# Change `s3` to `gs` for Google Cloud Storage

resp = bulk_import(
api_key=API_KEY,
url=CLOUD_API_ENDPOINT,
cluster_id=CLUSTER_ID,
collection_name=COLLECTION_NAME,
object_url=object_url,
access_key=ACCESS_KEY,
secret_key=SECRET_KEY
)

job_id = resp.json()['data']['jobId']
print(job_id)

# job-0103f039ccdq9aip1xd4rf
📘メモ

object_url は、リモート bucket 内のファイルまたはフォルダへの有効な URL である必要があります。提供されているコードでは、format() メソッドを使用して bucket 名と writer が返すデータパスを結合し、有効な object URL を作成しています。

データとターゲット collection の両方が AWS 上でホストされている場合、object URL は s3://remote-bucket/file-path のようになります。writer が返したデータパスの前に付ける URI については、Storage Options を参照してください。

タスクの進行状況を確認する

以下のコードは、5 秒ごとに bulk-import の進行状況を確認し、進捗率をパーセンテージで出力します。

python
import time
from pymilvus import get_import_progress

job_id = res.json()['data']['jobId']

res = get_import_progress(
api_key=API_KEY,
url=CLOUD_API_ENDPOINT,
cluster_id=CLUSTER_ID, # Zilliz Cloud cluster ID, like "in01-xxxxxxxxxxxxxxx"
job_id=job_id,
)

print(res.json()["data"]["progress"])

# check the bulk-import progress
while res.json()["data"]["progress"] < 100:
time.sleep(5)

res = get_import_progress(
url=CLOUD_API_ENDPOINT,
api_key=API_KEY,
job_id=job_id,
cluster_id=CLUSTER_ID
)

print(res.json()["data"]["progress"])

# 0 -- import progress 0%
# 49 -- import progress 49%
# 100 -- import finished
📘メモ

get_import_progress() 内の url は、ターゲット collection のクラウドリージョンに対応するものに置き換えてください。

以下のように、すべての bulk-import ジョブを一覧表示できます。

python
from pymilvus import list_import_jobs

res = list_import_jobs(
api_key=API_KEY,
url=CLOUD_API_ENDPOINT,
cluster_id=CLUSTER_ID # Zilliz Cloud cluster ID, like "in01-xxxxxxxxxxxxxxx"
)

print(res.json())

# {
# "code": 0,
# "data": {
# "records": [
# {
# "collectionName": "zero_to_hero",
# "jobId": "job-01f36d8fd67u94avjfnxi0",
# "state": "Completed"
# }
# ],
# "count": 1,
# "currentPage": 1,
# "pageSize": 10
# }
# }

まとめ

このコースでは、データをインポートする一連のプロセス全体を扱いました。最後に、重要なポイントを振り返ります。

  • データを確認して、ターゲット collection の schema を設計します。

  • BulkWriter を使用する際は、次の点に注意してください。

    • 追加する各行に、schema で定義されたすべてのフィールドをキーとして含めてください。dynamic fields が許可されている場合は、該当する未定義フィールドも含めてください。

    • すべての行を追加した後、commit() の呼び出しを忘れないでください。

  • bulk_import() を使用する際は、準備済みデータをホストするクラウドプロバイダの endpoint と、writer が返したデータパスを連結して object URL を構築します。

Ctrl I