OpenSearch SDK for Java V3.1 を使用して、コミットモードでドキュメントの追加、更新、削除を行います。コミットモードでは、クライアント側キューに 1 つ以上のドキュメントをバッファーし、単一の commit 呼び出しでバッチ全体を送信します。
ユースケース
複数のソースから動的にデータをマージしてコミットする
リアルタイムで単一のドキュメントをコミットする
小規模なドキュメントバッチを一度にコミットする
コミットモードでは、同一テーブル内のフィールドのみがサポートされます。異なるテーブルのフィールドを 1 回の呼び出しでコミットすることはできません。
前提条件
開始する前に、以下の準備を完了していることを確認してください。
構成済みのテーブルを持つ OpenSearch アプリケーション
ご利用のリージョンに対応するアプリケーション名 (
appName)、テーブル名 (tableName)、および OpenSearch API エンドポイント (host)必要な OpenSearch 権限を持つ Resource Access Management (RAM) ユーザーからの AccessKey ID および AccessKey Secret。設定手順については、「RAM ユーザーの作成」および「アクセス権限付与ルール」をご参照ください。
お使いの Alibaba Cloud アカウントの AccessKey ペアは、すべての API オペレーションにアクセスできます。代わりに RAM ユーザーの AccessKey ペアを使用し、その RAM ユーザーに AliyunServiceRoleForOpenSearch ロールを付与してください。認証情報をソースコード内に直接保存しないでください。
環境変数の設定
サンプルコードが実行時に認証情報を読み取れるよう、環境変数として設定します。
Linux および macOS
<access_key_id> および <access_key_secret> を RAM ユーザーの AccessKey ID および AccessKey Secret に置き換え、以下のコマンドを実行します。
export ALIBABA_CLOUD_ACCESS_KEY_ID=<access_key_id>
export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<access_key_secret>Windows
環境変数ファイルを作成し、
ALIBABA_CLOUD_ACCESS_KEY_IDおよびALIBABA_CLOUD_ACCESS_KEY_SECRETを対応する値とともに追加します。変更を有効にするために Windows を再起動します。
仕組み
各操作は、以下の 3 ステップのパターンに従います。
ドキュメントデータを
Mapオブジェクトにカプセル化します。add()、update()、またはremove()を呼び出して、ドキュメントをクライアント側バッファーにステージングします。commit(appName, tableName)を呼び出して、バッファー内のすべてのドキュメントを OpenSearch に送信します。
コミットが正常に完了した後は、OpenSearch コンソールのエラーログを確認し、フィールドレベルのエラー(型変換失敗など)が発生していないかを検証してください。
ドキュメントの追加
ドキュメントのフィールドを含む Map オブジェクトを作成し、add() でステージングした後、commit() で送信します。
package com.aliyun.opensearch;
import com.aliyun.opensearch.sdk.dependencies.com.google.common.collect.Maps;
import com.aliyun.opensearch.sdk.generated.OpenSearch;
import com.aliyun.opensearch.sdk.generated.commons.OpenSearchClientException;
import com.aliyun.opensearch.sdk.generated.commons.OpenSearchException;
import com.aliyun.opensearch.sdk.generated.commons.OpenSearchResult;
import java.io.UnsupportedEncodingException;
import java.util.Map;
import java.util.Random;
public class testCommitSearch {
private static String appName = "データをコミットする OpenSearch アプリケーションの名称";
private static String tableName = "データをアップロードするテーブルの名称";
private static String host = "ご利用のリージョンにおける OpenSearch API のエンドポイント";
public static void main(String[] args) {
// 環境変数から認証情報を読み取ります
String accesskey = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID");
String secret = System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
// デバッグ用:ファイルエンコーディングとデフォルト文字セットを出力
System.out.println(String.format("file.encoding: %s", System.getProperty("file.encoding")));
System.out.println(String.format("defaultCharset: %s", java.nio.charset.Charset.defaultCharset().name()));
// ドキュメントのプライマリキーとしてランダム整数を使用
Random rand = new Random();
int value = rand.nextInt(Integer.MAX_VALUE);
// Map オブジェクトとしてドキュメントを構築
Map<String, Object> doc1 = Maps.newLinkedHashMap();
doc1.put("id", value);
String title_string = "コミットモードで doc1 をアップロード"; // UTF-8
byte[] bytes;
try {
bytes = title_string.getBytes("utf-8");
String utf8_string = new String(bytes, "utf-8");
doc1.put("name", utf8_string);
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
doc1.put("phone", "1381111****");
int[] int_arr = {33, 44};
doc1.put("int_arr", int_arr);
String[] literal_arr = {"コミットモードで doc1 をアップロード", "コミットモードでの doc1 アップロードをテスト"};
doc1.put("literal_arr", literal_arr);
float[] float_arr = {(float) 1.1, (float) 1.2};
doc1.put("float_arr", float_arr);
doc1.put("cate_id", 1);
// SDK クライアントチェーンを初期化:OpenSearch → OpenSearchClient → DocumentClient
OpenSearch openSearch1 = new OpenSearch(accesskey, secret, host);
OpenSearchClient serviceClient1 = new OpenSearchClient(openSearch1);
DocumentClient documentClient1 = new DocumentClient(serviceClient1);
// ドキュメントをクライアント側バッファーにステージング
documentClient1.add(doc1);
System.out.println(doc1.toString());
try {
// バッファー内のすべてのドキュメントを送信します。commit() の呼び出し前に複数のドキュメントをバッファーできます。
OpenSearchResult osr = documentClient1.commit(appName, tableName);
checkCommitResult(osr);
} catch (OpenSearchException | OpenSearchClientException e) {
e.printStackTrace();
}
// OpenSearch は非同期でデータをインデックス化するため、クエリ実行前に 10 秒待機します
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
// コミット結果をチェックします。トランスポートレベルのエラーがない場合に true が返されます。
// 「true」が返された場合、リクエストは正常に受信されたことを意味します。
// フィールドレベルのエラー(例:型変換失敗)は、OpenSearch コンソールのエラーログで別途報告されます。
private static void checkCommitResult(OpenSearchResult osr) {
if (osr.getResult().equalsIgnoreCase("true")) {
System.out.println("データのコミット中にエラーは発生しませんでした!\n リクエスト ID は " + osr.getTraceInfo().getRequestId());
} else {
System.out.println("データのコミット中にエラーが発生しました!" + osr.getTraceInfo());
}
}
}ドキュメントの更新
ドキュメントを更新するには、元のドキュメントと同じプライマリキーと更新後のフィールド値を含む新しい Map を構築し、update() を呼び出した後、commit() を実行します。
// 更新後のドキュメントを構築 — 元のドキュメントと同じプライマリキーを含める必要があります
Map<String, Object> doc2 = Maps.newLinkedHashMap();
doc2.put("id", value); // doc1 と同じプライマリキー
String title_string2 = "コミットモードで doc1 を更新"; // UTF-8
byte[] bytes2;
try {
bytes2 = title_string2.getBytes("utf-8");
String utf8_string2 = new String(bytes2, "utf-8");
doc2.put("name", utf8_string2);
} catch (UnsupportedEncodingException e) {
e.printStackTrace();
}
doc2.put("phone", "1390000****");
int[] int_arr2 = {22, 22};
doc2.put("int_arr", int_arr2);
String[] literal_arr2 = {"コミットモードで doc1 を更新", "コミットモードで doc1 を更新"};
doc2.put("literal_arr", literal_arr2);
float[] float_arr2 = {(float) 1.1, (float) 1.2};
doc2.put("float_arr", float_arr2);
doc2.put("cate_id", 1);
// 更新をクライアント側バッファーにステージング
documentClient1.update(doc2);
System.out.println(doc2.toString());
try {
OpenSearchResult osr = documentClient1.commit(appName, tableName);
checkCommitResult(osr);
} catch (OpenSearchException | OpenSearchClientException e) {
e.printStackTrace();
}
// クエリ実行前に 10 秒待機
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
e.printStackTrace();
}ドキュメントの削除
ドキュメントを削除するには、プライマリキーのみを含む Map を構築し、remove() を呼び出した後、commit() を実行します。
// 削除対象ドキュメントを特定するには、プライマリキーのみが必要です
Map<String, Object> doc3 = Maps.newLinkedHashMap();
doc3.put("id", value);
// 削除をクライアント側バッファーにステージング
documentClient1.remove(doc3);
System.out.println(doc3.toString());
try {
OpenSearchResult osr = documentClient1.commit(appName, tableName);
checkCommitResult(osr);
} catch (OpenSearchException | OpenSearchClientException e) {
e.printStackTrace();
}
// クエリ実行前に少なくとも 1 秒(推奨:10 秒)待機します。
// 削除直後にクエリを実行すると、削除済みのドキュメントがまだ返される可能性があります。
try {
Thread.sleep(10000);
} catch (InterruptedException e) {
e.printStackTrace();
}結果の確認
各コミット後に OpenSearch に対してクエリを実行し、期待通りにデータが存在するか、または存在しないかを確認します。
import com.aliyun.opensearch.sdk.dependencies.com.google.common.collect.Lists;
import com.aliyun.opensearch.sdk.dependencies.org.json.JSONObject;
import com.aliyun.opensearch.sdk.generated.search.Config;
import com.aliyun.opensearch.sdk.generated.search.SearchFormat;
import com.aliyun.opensearch.sdk.generated.search.SearchParams;
import com.aliyun.opensearch.sdk.generated.search.general.SearchResult;
// 検索専用の別クライアントを初期化
OpenSearch openSearch2 = new OpenSearch(accesskey, secret, host);
OpenSearchClient serviceClient2 = new OpenSearchClient(openSearch2);
SearcherClient searcherClient2 = new SearcherClient(serviceClient2);
// ページングおよび応答フォーマットを設定
// サポートされるフォーマット:XML、JSON。FULLJSON はサポートされていません。
Config config = new Config(Lists.newArrayList(appName));
config.setStart(0);
config.setHits(30);
config.setSearchFormat(SearchFormat.JSON);
// プライマリキーによるクエリ
SearchParams searchParams = new SearchParams(config);
searchParams.setQuery("id:'" + value + "'");
try {
SearchResult searchResult = searcherClient2.execute(searchParams);
String result = searchResult.getResult();
JSONObject obj = new JSONObject(result);
System.out.println("クエリ結果: " + obj.toString());
} catch (OpenSearchException | OpenSearchClientException e) {
e.printStackTrace();
}次のステップ
AliyunServiceRoleForOpenSearch -- サービスリンクロールに必要な権限を確認する
アクセス承認ルール-RAM ユーザー向けの詳細なアクセスを設定
AccessKey ペアの作成 -- RAM ユーザー向けの認証情報を生成します