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

Spark バッチジョブ

この機能は AWS us-west-2 リージョンでのみ利用できます。Google Cloud および Microsoft Azure では利用できません。

Spark バッチジョブを使用すると、Zilliz Cloud で管理されている大規模データセットに対して分散オフライン処理を実行できます。組み込みジョブを使用して、ベクトルデータの重複排除、クラスター化、検査を行えます。

Spark バッチジョブは、長時間実行されるデータ処理タスク向けに設計されています。低レイテンシーのオンラインリクエストやレコード単位の変換を目的としたものではありません。

発生する可能性のある問題

ベクトルデータセットが大きくなるにつれて、単純な挿入や検索操作だけでは対応できなくなる場合があります。繰り返しの取り込みによって重複が発生したり、大規模な埋め込みコレクションの全体像を把握しにくくなったり、前処理パイプラインの失敗によって疑わしいレコードが残ったりすることがあります。

重複するベクトル埋め込み

再試行、繰り返しのインポート、データソースの重複、同じテキストや画像のわずかな変更版などによって、エンティティ間で重複するベクトル埋め込みが作成されることがあります。一部のエンティティは同じプライマリキーを共有していますが、ほかのエンティティはプライマリキーが異なるものの内容がほぼ同一であり、ストレージ、インデックス作成、後続処理のコストを増加させます。

Spark バッチジョブは、大規模なデータセット全体で重複を特定して削除できます。同じ ID を持つレコードを整理するには プライマリキー重複排除ジョブ を、ID は異なるが内容がほぼ同一のレコードを検出するには ベクトル類似度重複排除ジョブ を使用します。

不明瞭な埋め込み分布

コレクションが大きくなるにつれて、埋め込み分布の全体像を把握することが難しくなります。どのパターンがデータセットを支配しているのか、ロングテールデータがどこに現れるのか、新しくインポートされたデータが既存のレコードと異なるのかどうかを把握できない場合があります。

K-Means クラスタリングジョブ は、類似する埋め込みを大まかなクラスターにグループ化し、各レコードにクラスター ID を割り当てます。この結果を利用して、データ分布の分析、データソースの比較、代表的なサンプルの作成、類似度に基づく処理の小規模なグループへの分割を行えます。

埋め込みデータに潜む異常

埋め込みパイプラインでは、一見すると有効に見えるものの、データセットの残りの部分と一致しないレコードが生成されることがあります。前処理の失敗、不適切なモデル、パースノイズ、破損したソースコンテンツ、予期しないデータバッチなどによって、手動での確認では発見しにくい異常な埋め込みが作成される可能性があります。

異常検知ジョブ は、埋め込み分布をスキャンし、一般的なパターンから大きく外れるレコードを特定します。この結果を利用して、再埋め込み、クリーンアップ、またはさらなるレビューが必要になる可能性があるデータを見つけることができます。異常は必ずしも無効であるとは限らないため、フラグが付けられたレコードは自動的に削除するのではなく、レビューする必要があります。

ジョブタイプの選択

次の表に、目的に応じて推奨されるジョブタイプを示します。

目的推奨されるジョブ
同じプライマリキーを共有するレコードの削除プライマリキー重複排除
ベクトル表現が非常に類似しているレコードの検出ベクトル類似度重複排除
ベクトルデータをあらかじめ定義された数のグループに分割K-Means クラスタリング
主要なデータ分布から大きく外れるレコードの検出異常検知

Spark バッチジョブの仕組み

Spark バッチジョブは、長時間実行される分散オフライン処理ジョブであり、ジョブ作成リクエストを受け取るとすぐにジョブ ID を返します。このジョブ ID をハンドラーとして使用して、進行状況の監視とライフサイクルの管理を行えます。

事前準備

Spark バッチジョブを送信するには、次の条件を満たしていることを確認してください。

  • 十分な権限を持つ有効な Zilliz Cloud API キーがあること。

  • 入力データが、サポートされている形式で Zilliz Cloud ボリュームに用意されていること。

    • サポートされているデータファイル形式は parquetlancejsoncsv です。
  • 各外部ボリュームに関連付けられたストレージロールが、ジョブに必要な権限を持っていること。

    • 入力用の外部ボリュームには、入力場所への読み取りアクセス権が必要です。

    • 出力用の外部ボリュームには、出力場所への読み取りと書き込みのアクセス権が必要です。これには、ジョブの実行中に作成された一時オブジェクトや不完全なオブジェクトを削除する権限も含まれます。

    • 入力と出力の外部ボリュームでは、異なるストレージロールを使用できます。

外部ボリュームの権限を構成する

Spark バッチジョブは、外部ボリュームから入力データを読み取り、結果を外部ボリュームに書き込みます。ジョブを送信する前に、各ボリュームに関連付けられたオブジェクトストレージのロールが、その用途に必要な権限を持っていることを確認してください。

入力ボリュームと出力ボリュームでは、異なるストレージロールを使用できます。各ロールには、対応するバケットとプレフィックスに必要な権限のみを付与してください。

入力ボリューム

入力ボリュームには、入力データへの読み取りアクセス権が必要です。ボリュームで使用する統合がすでに必要な読み取り権限を提供している場合、追加の権限は必要ありません。

Amazon S3 の場合、ロールには次の権限が必要です。

  • 入力プレフィックス配下のオブジェクトに対する s3:GetObject

  • 入力プレフィックス配下のオブジェクトを一覧表示するための s3:ListBucket

  • バケットに対する s3:GetBucketLocation

例を表示するには、ここをクリックしてください。
json
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "ListInput",
"Effect": "Allow",
"Action": [
"s3:ListBucket",
"s3:GetBucketLocation"
],
"Resource": "arn:aws:s3:::<bucket>",
"Condition": {
"StringLike": {
"s3:prefix": [
"<input-prefix>",
"<input-prefix>/*"
]
}
}
},
{
"Sid": "ReadInput",
"Effect": "Allow",
"Action": "s3:GetObject",
"Resource": "arn:aws:s3:::<bucket>/<input-prefix>/*"
}
]
}

出力ボリューム

出力ボリュームには、ジョブの結果を書き込む権限が必要です。また、Spark は一時オブジェクトや、失敗またはキャンセルされた書き込みによって残されたオブジェクトを削除する場合があります。

Amazon S3 の場合、ロールには次の権限が必要です。

  • 出力場所にアクセスするための s3:GetObjects3:ListBuckets3:GetBucketLocation

  • 結果を書き込むための s3:PutObject

  • 一時的な出力オブジェクトや不完全な出力オブジェクトをクリーンアップするための s3:DeleteObject

s3:PutObjects3:DeleteObject は、バケット全体ではなく出力プレフィックスに限定してください。

例を表示するには、ここをクリックしてください。
json
{
"Version": "2012-10-17",
"Statement": [
{
"Sid": "ListOutput",
"Effect": "Allow",
"Action": [
"s3:ListBucket",
"s3:GetBucketLocation"
],
"Resource": "arn:aws:s3:::<bucket>",
"Condition": {
"StringLike": {
"s3:prefix": [
"<output-prefix>",
"<output-prefix>/*"
]
}
}
},
{
"Sid": "ReadWriteOutput",
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:DeleteObject"
],
"Resource": "arn:aws:s3:::<bucket>/<output-prefix>/*"
}
]
}

Spark バッチジョブを送信して実行する

API キーを取得し、必要なファイルを Zilliz Cloud ボリュームにアップロードしたら、Spark バッチジョブを送信して実行する準備は完了です。

1

べき等性キーを生成します。

べき等性キーは、同じジョブリクエストを再試行するときに変わらない一意の文字列です。詳細については、べき等性のある送信 を参照してください。

2

リクエストヘッダーを準備します。

Spark バッチジョブを作成するときは、前のステップで生成したべき等性キーをリクエストヘッダーに含めます。

http
Authorization: Bearer <api-key>
Idempotency-Key: spark-job-20260730-001
Content-Type: application/json
3

リクエストペイロードを準備します。

すべての Spark バッチジョブは共通のペイロード構造を共有しており、ジョブの目的に応じてジョブ固有のパラメーターが異なります。

共通のペイロード構造については、リクエストペイロード を参照してください。ジョブ固有のパラメーターについては、次のページを参照してください。

べき等性のある送信

べき等性キーを使用すると、ジョブの送信を安全に再試行できます。Zilliz Cloud はキーとリクエストボディの両方を確認して、既存のジョブを返すか、リクエストを競合として拒否するかを判断します。一致時の動作を次の表に示します。

ケース動作
べき等性キーが同じでリクエストボディも同じ以前に作成された Spark バッチジョブを返します。重複するジョブは作成されません。
べき等性キーが同じでリクエストボディが異なる競合エラーを返します。

べき等性キーのスコープは、特定の組織、ユーザー、リージョン、およびプロジェクトです。Zilliz Cloud は、ジョブの最大タイムアウト期間に基づいて、キーを約 25 時間保持します。保持期間が過ぎると、同じキーを新しい送信に再利用できます。

リクエストペイロード

すべての Spark バッチジョブは、次のような共通のペイロード構造を共有します。

json
{
"description": "optional description",
"regionId": "aws-us-west-2",
"input": {...},
"output": {...},
"resourceSize": "SMALL",
"timeoutSeconds": 3600
}

次の表に、これらのパラメーターの説明を示します。

パラメーター必須説明
descriptionいいえジョブの説明(任意)です。
値は 1,024 文字以内の文字列です。
regionIdはいジョブを実行する Zilliz Cloud リージョンの ID です。サポートされているリージョンについては、Cloud Providers & Regions を参照してください。
inputいいえジョブの入力です。これは組み込みジョブでのみ必須です。詳細については、以下の表を参照してください。
outputいいえジョブの出力です。これは組み込みジョブで必須です。詳細については、以下の表を参照してください。
resourceSizeいいえジョブに必要な Spark クラスターのサイズです。指定可能な値は SMALLMEDIUMLARGEXLARGE2XLARGE3XLARGE です。
timeoutSecondsいいえ現在のジョブのタイムアウト期間(秒)です。値は 300 から 86400 までの正の整数です。

上記の表の input パラメーターと output パラメーターは、次のような同様の構造を共有します。

パラメーター

必須

説明

type

はい

Spark バッチジョブのタイプです。これは inputoutput の両方に適用されます。指定可能な値は次のとおりです。

  • volume

volumeName

いいえ

Zilliz Cloud ボリュームの名前です。typevolume に設定する場合に必須です。これは inputoutput の両方に適用されます。

path

いいえ

input/output ファイルのパスは、指定した Zilliz Cloud ボリュームのルートからの相対パスです。typevolume に設定する場合に必須です。これは inputoutput の両方に適用されます。

ファイルが volume://path/to/data.parquet にある場合は、pathpath/to/data.parquet に設定します。

format

いいえ

入力ファイルまたは出力ファイルの形式です。これは inputoutput の両方に適用されます。このパラメーターのデフォルト値は parquet です。指定可能な値は parquetlancejsoncsv です。

writeMode

いいえ

出力ファイルの書き込みモードです。これは output にのみ適用されます。指定可能な値は次のとおりです。

  • ERROR_IF_EXIST

    指定した出力ファイルがすでに存在する場合にエラーを返します。これがデフォルトのオプションです。

  • OVERWRITE

    指定したファイルを上書きします。

Notes

input の場合、format によって使用される Spark データソースリーダーが決まります。このパラメーターを省略すると、ジョブはデフォルトで Parquet リーダーを使用し、指定したパス配下の Parquet ファイルのみを処理します。JSON や CSV などほかの形式のファイルは無視されます。

処理するファイルが Parquet 以外の場合は、input.format に対応する形式を明示的に設定してください。このジョブは、異なる形式のファイルを自動的に検出して結合することはありません。

送信レスポンス

成功したジョブリクエストのレスポンスでは、異なる HTTP コードが返される場合があります。次の表に、レスポンスに含まれる該当する HTTP コードを示します。

ケースHTTP コード説明
ジョブの作成201 CREATEDジョブが作成中であることを示します。
ジョブの作成は非同期であるため、レスポンスに含まれるジョブ ID を使用して、進行状況を確認したりライフサイクルを管理したりできます。
ジョブのキャンセル202 ACCEPTEDキャンセルが受け付けられ、処理中であることを示します。
ジョブのキャンセルは非同期であるため、レスポンスに含まれるジョブ ID を使用して進行状況を確認できます。
ジョブの詳細取得200 OKリクエストされたレスポンスが返されたことを示します。
これは同期処理であり、レスポンスにはリクエストが処理された時点のジョブステータスが常に含まれます。

HTTP コードは異なりますが、次のように同じペイロード構造を共有します。

json
{
"code": 0,
"data": {
"jobId": "job-xxxxxxxx",
"projectId": "proj-xxxxxxxx",
"type": "SPARK",
"description": "backfill product attributes",
"status": "PENDING",
"regionId": "aws-us-west-2",
"clusterId": "in-xxxxxxxx",
"createdAt": null,
"startedAt": null,
"finishedAt": null,
"durationSeconds": null
}
}

レスポンスには、ジョブの監視に使用できるジョブ ID が含まれます。ジョブの状態を確認したり、ジョブの詳細を取得したり、実行中のジョブをキャンセルしたりするには、ジョブの状態を理解する を参照してください。

エラーが発生した場合、レスポンスは次のようになります。

json
{
"code": 10001,
"message": "projectId is required",
"details": {
"errorCode": "INVALID_PARAMETER"
}
}

エラーの内容は、details.errorCode と、エラーレスポンスに含まれる HTTP コードから判断できます。次の表に、該当する HTTP コードとその意味を示します。

HTTP コード説明
400 BAD REQUESTリクエストに含まれるパラメーターが正しくないことを示します。
403 FORBIDDENAPI キーに十分な権限がないか、指定したプロジェクトでリソースを利用できないことを示します。
404 NOT FOUNDジョブ ID や Zilliz Cloud ボリュームなど、指定したリソースが存在しないことを示します。
409 CONFLICTべき等性キーがリクエストペイロードと一致しないことを示します。
500 INTERNAL SERVER ERRORサーバーがリクエストの処理に失敗したことを示します。

次のステップ

次のガイドを参照して、目的に合った Spark バッチジョブを作成するか、既存のジョブの監視と管理を行ってください。

プライマリキー重複排除

プライマリキー重複排除は、同じプライマリキーを共有するレコードを特定し、大規模なデータセットから冗長なコピーを削除します。このジョブを使用すると、繰り返しのインポート、パイプラインの再試行、移行、または重複するデータソースによって発生した重複をクリーンアップできます。 | Cloud

ベクトル類似度重複排除

ベクトル類似度重複排除は、非常に類似した埋め込みを持つレコードを特定し、それらを意味的な重複としてグループ化します。このジョブを使用すると、言い換えられたテキスト、わずかに変更された画像、類似コンテンツの複数バージョンなど、意味的な冗長性を削減できます。 | Cloud

K-Means クラスタリング

K-Means クラスタリングは、類似した埋め込みを持つレコードを指定された数のクラスターにグループ化します。このジョブを使用すると、ベクトルデータの分布を把握したり、レコードを大まかな意味グループに整理したり、サンプリングや分析などの後続ワークフロー向けにデータセットを準備できます。 | Cloud

異常検知

異常検知は、ベクトル埋め込みがデータ全体の分布から大きく逸脱しているレコードを特定します。このジョブを使用すると、データ品質の問題、まれなケース、処理エラー、または追加の確認が必要なサンプルを示唆する可能性がある異常なレコードを見つけられます。 | Cloud

データバックフィル

データバックフィルを使用すると、Zilliz Cloud コレクション内の既存エンティティの選択したフィールドを、Zilliz Cloud ボリューム上の Parquet ファイルに保存されたデータを使って更新できます。バックフィルジョブは、入力レコードを主キーによって既存エンティティと照合し、指定したフィールド値をコレクションに書き戻します。新しく追加したフィールドへの値の設定、欠損値の補完、既存フィールド値の一括置換などに使用できます。 | Cloud

Spark バッチジョブの管理

Spark バッチジョブは非同期で実行され、送信から完了までの間に複数の状態を遷移します。このページでは、ジョブのライフサイクルについて説明したうえで、ジョブの一覧表示、ジョブの詳細の取得、およびキャンセル可能な状態にあるジョブをキャンセルする方法を説明します。 | Cloud