すべてのプロダクト
Search
ドキュメントセンター

Platform For AI:Java SDK

最終更新日:Aug 29, 2026

このガイドでは、Java SDK を使用して Elastic Algorithm Service (EAS) のモデルサービスを呼び出す方法を、入出力の例やサンプルプログラムとあわせて説明します。

説明

SDK のユースケースと原則については、「サービス呼び出し SDK」をご参照ください。

前提条件

Maven プロジェクトで EAS Java SDK を使用するには、pom.xml ファイルの <dependencies> セクションに eas-sdk の依存関係を追加してください。最新バージョンについては、Maven リポジトリ を確認してください。

<dependency>
  <groupId>com.aliyun.openservices.eas</groupId>
  <artifactId>eas-sdk</artifactId>
  <version>2.0.20</version>
</dependency>

EAS SDK 2.0.5 以降には、マルチ優先度非同期キューサービス向けの QueueService クライアント機能が含まれています。この機能を使用し、依存関係の競合を回避するには、次の 2 つの依存関係を追加し、必要に応じてそれぞれのバージョンを調整してください:

<dependency>
    <groupId>org.java-websocket</groupId>
    <artifactId>Java-WebSocket</artifactId>
    <version>1.5.1</version>
</dependency>
<dependency>
    <groupId>org.apache.commons</groupId>
    <artifactId>commons-lang3</artifactId>
    <version>3.1</version>
</dependency>

クイックスタート

Java SDK を使用してサービス呼び出しを行うには、次の 3 つのステップを実行します。

  1. 呼び出し情報の取得:EAS コンソールのサービス詳細ページで、[Call Information] タブに移動し、エンドポイント、サービス名、トークンを取得します。

  2. リクエストタイプの選択とコードの記述:モデルの入力データ形式に基づいて適切なリクエスト/レスポンス クラスを選択し、以下の最小限の例を使用してコードを記述します。

    説明

    組み込みプロセッサを使用してサービスをデプロイした場合、SDK は対応する入出力クラスを提供します。たとえば、組み込みのTensorFlow プロセッサTFRequest に対応します。詳細については、組み込みプロセッサの各プロセッサのドキュメントをご参照ください。

  3. 実行と検証:クライアントプログラムを実行し、レスポンスを検証します。エラーが発生した場合は、トラブルシューティング ガイドをご参照ください。

次のコードは、文字列リクエストの最小限のエンドツーエンドの例です。その他の例については、「プログラム例」をご参照ください。

import com.aliyun.openservices.eas.predict.http.PredictClient;
import com.aliyun.openservices.eas.predict.http.HttpConfig;

public class TestString {
    public static void main(String[] args) throws Exception {
        PredictClient client = new PredictClient(new HttpConfig());
        
        // VPC 専用接続を使用するには、setDirectEndpoint メソッドを呼び出します。形式は通常 {uid}.vpc.{region-id}.pai-eas.aliyuncs.com です。
        client.setDirectEndpoint("182848887922****.vpc.cn-shanghai.aliyuncs.com");
        // EAS サービスのパブリックエンドポイント。形式は通常 {uid}.{region-id}.pai-eas.aliyuncs.com です。
        // client.setEndpoint("182848887922****.cn-shanghai.pai-eas.aliyuncs.com");
        
        // EAS サービスの名前。
        client.setModelName("your_service_name");
        client.setToken("YOUR_SERVICE_TOKEN");
        // リクエストパス。完全なリクエスト URL は http://<endpoint>/api/direct/<modelName>/<requestPath> です。
        client.setRequestPath("your_custom_path");
        // リクエストボディを構築します。サポートされる入力クラスは SDK によって異なります。この例では文字列を使用します。
        String request = "[{}]";
        String response = client.predict(request);
        System.out.println(response);

        client.shutdown();
    }
}

API リファレンス

Java SDK には、次のクラスが用意されています。

グループ

クラスの説明

メインクライアントクラス

PredictClient :エンドポイントやトークンなどのサービス詳細の設定、リクエストの送信、レスポンスの受信を行うためのメインクラスです。

接続設定

HttpConfig :タイムアウトや最大接続数など、HTTP 接続パラメータを設定します。

入出力

  • TFRequest :TensorFlow モデルへのリクエストをカプセル化します。

  • TFResponse :TensorFlow モデルからのレスポンスを解析します。

  • 文字列ベースのシナリオ では、専用のリクエストクラスまたはレスポンスクラスは不要です。 String オブジェクトとしてデータを直接受け渡しできます。

  • サポートされているその他の型の詳細については、SDK をご参照ください。

キューサービス

  • QueueClient :データの送信やサブスクライブを行う非同期キュークライアントです。このクラスを使用するには、前提条件 に記載されているとおり、追加の依存関係が必要です。

  • DataFrame :キューサービスのデータ項目をカプセル化するオブジェクトです。

PredictClient クラス

メインクライアントクラスです。サービス情報を設定し、リクエストを送信し、予測結果を受信します。

API

説明

PredictClient(HttpConfig httpConfig)

  • PredictClient インスタンスを作成します。

  • パラメータhttpConfigHttpConfig クラスのインスタンスです。

void setToken(String token)

  • HTTP リクエストの認証トークンを設定します。

  • パラメータtoken はサービスアクセス用の認証トークンです。

void setModelName(String modelName)

  • オンライン予測サービスのモデル名を設定します。

  • パラメータmodelName は使用するモデルの名前です。

void setEndpoint(String endpoint)

  • リクエストするサービスのホストとポートを指定します。形式は "host:port" です。

  • パラメータendpoint"host:port" 形式のサービスエンドポイントです。

void setDirectEndpoint(String endpoint)

  • VPC 専用接続経由でサービスにアクセスするためのエンドポイントを設定します。

  • パラメータendpoint はサービスエンドポイントです。

void setRequestPath(String requestPath)

  • サーバーサイドコードで定義されたリクエストパスを設定します。

  • パラメータrequestPath はサーバーサイドのリクエストパスです。例:client.setRequestPath("/custom_path")

void setRetryCount(int retryCount)

  • 失敗したリクエストのリトライ回数を設定します。

  • パラメータretryCount はリトライ回数です。

void setRetryConditions(EnumSet retryConditions)

  • リクエストをリトライする条件を設定します。このメソッドは setRetryCount メソッドと組み合わせて使用できます。デフォルトでは、すべてのリクエストエラーがリトライされます。このメソッドを使用すると、特定のリクエストエラーのみをリトライするように指定できます。

  • パラメータretryConditions は 1 つ以上のリトライ条件を含む EnumSet です。サポートされる条件は次のとおりです。

    • RetryCondition.CONNECTION_FAILED:リクエスト接続が失敗しました。

    • RetryCondition.CONNECTION_TIMEOUT:リクエスト接続がタイムアウトしました。

    • RetryCondition.READ_TIMEOUT:レスポンスの待機中にリクエストがタイムアウトしました。

    • RetryCondition.RESPONSE_5XX:サーバーが 5xx ステータスコードを返しました。

    • RetryCondition.RESPONSE_4XX:サーバーが 4xx ステータスコードを返しました。

  • client.setRetryConditions(
        EnumSet.of(
            RetryCondition.READ_TIMEOUT,    // 読み取りタイムアウト時にリトライ
            RetryCondition.RESPONSE_5XX     // 5xx エラーコード時にリトライ
        )
    );

    この例では、リクエストがタイムアウトした場合、またはサーバーが 5xx ステータスコードを返した場合にのみリトライするように指定しています。

void setContentType(String contentType)

  • HTTP リクエストの Content-Type を設定します。デフォルトは "application/octet-stream" です。

  • パラメータcontentType は送信するデータストリームのコンテンツタイプです。

void setUrl(String url)

カスタムのリクエスト URL を設定します。

void setCompressor(Compressor compressor)

  • リクエストデータの圧縮方法を設定します。

  • パラメータcompressor は圧縮方法です。サポートされる値は Compressor.GzipCompressor.Zlib です。

  • 詳細については、「リクエストデータ圧縮の例」をご参照ください。

void addExtraHeaders(Map<String, String> extraHeaders)

  • リクエストにカスタム HTTP ヘッダーを追加します。

  • パラメータextraHeaders は追加する HTTP ヘッダーの Map<String, String> です。

PredictClient createChildClient(String token, String endpoint, String modelName)

  • 親クライアントのスレッドプールを共有する子クライアントを作成して返します。これはマルチスレッド予測に役立ちます。

  • パラメータ

    • token:サービスの認証トークンです。

    • endpoint:サービスのエンドポイントです。

    • modelName:モデルの名前です。

void predict(TFRequest runRequest)

  • TensorFlow リクエストをオンライン予測サービスに送信します。

  • パラメータrunRequestTFRequest のインスタンスです。

void predict(String requestContent)

  • 文字列ベースのリクエストをオンライン予測サービスに送信します。

  • パラメータrequestContent は文字列形式のリクエスト内容です。

void predict(byte[] requestContent)

  • バイト配列のリクエストをオンライン予測サービスに送信します。

  • パラメータrequestContent はバイト配列形式のリクエスト内容です。

HttpConfig クラス

タイムアウト、スレッド数、接続プールなど、基盤となる HTTP 接続パラメーターを設定します。

API

説明

void setIoThreadNum(int ioThreadNum)

  • HTTP リクエストの I/O スレッド数を設定します。デフォルトは 2 です。

  • パラメーター: ioThreadNum は、I/O スレッドの数です。

void setReadTimeout(int readTimeout)

  • 接続確立後、サーバーからのデータパケットを待機する最大時間です。デフォルトは 5000 (5 秒) です。

  • パラメーター: readTimeout はミリ秒単位の読み取りタイムアウトです。

重要

このタイムアウトは接続確立後にのみ適用されます。 setRequestTimeout メソッドで設定されるリクエストタイムアウトとは異なります。

void setRequestTimeout(int requestTimeout)

  • リクエストの送信から完全なレスポンスを受信するまでの合計時間です。デフォルトは 5000 (5 秒) です。

  • パラメーターrequestTimeout はリクエストタイムアウト (ミリ秒単位) です。

重要

このタイムアウトは、接続確立、データ転送、サーバー処理を含むリクエストのライフサイクル全体をカバーします。 setReadTimeout メソッドで設定される読み取りタイムアウトとは異なります。

void setConnectTimeout(int connectTimeout)

  • 接続確立時に待機する最大時間です。デフォルトは 5000 (5 秒) です。

  • パラメーター: connectTimeout は、ミリ秒単位の接続タイムアウトです。

void setMaxConnectionCount(int maxConnectionCount)

  • 接続プール内の最大合計接続数を設定します。デフォルトは 1000 です。

  • パラメーター: maxConnectionCount は最大接続数です。

void setMaxConnectionPerRoute(int maxConnectionPerRoute)

  • ルートごとの最大接続数を設定します。デフォルトは 1000 です。

  • パラメーター: maxConnectionPerRoute は各ルートの最大接続数です。

void setKeepAlive(boolean keepAlive)

  • 機能: HTTP サービスの keep-alive を設定します。

  • パラメーター: keepAlive。接続に対して keep-alive メカニズムを有効にするかどうかを指定します。デフォルト値は true です。

int getErrorCode()

直前の API 呼び出しのステータスコードを返します。

String getErrorMessage()

直前の API 呼び出しのステータスメッセージを返します。

TFRequest クラス

TensorFlow モデルの入力データを構築します。

API

説明

void setSignatureName(String value)

  • 機能: モデルが TensorFlow SavedModel フォーマットの場合、リクエストするモデルの signatureDef の名前を指定します。

  • パラメーター: リクエストするモデルの signatureDef の名前です。

void addFetch(String value)

  • モデルからフェッチする出力テンソルを指定します。

  • パラメーターvalue はフェッチする出力テンソルのエイリアスです。

void addFeed(String inputName, TFDataType dataType, long[] shape, ?[] content)

  • リクエストに入力テンソルを追加します。

  • パラメーター

    • inputName: 入力テンソルのエイリアスです。

    • dataType: 入力テンソルのデータ型です。

    • shape: 入力テンソルの形状です。

    • content: フラット化された 1 次元配列形式のテンソルの内容です。配列内の要素型は dataType に依存します。

      入力テンソルの DataType が DT_FLOAT、DT_COMPLEX64、DT_BFLOAT16、または DT_HALF の場合、content の要素型は FLOAT です。DataType が DT_COMPLEX64 の場合、content 内の隣接する 2 つの FLOAT 要素は、それぞれ複素数の実部と虚部を表します。

      入力テンソルの DataType が DT_DOUBLE または DT_COMPLEX128 の場合、content の要素は DOUBLE 型です。DataType が DT_COMPLEX128 の場合、content 内の隣接する 2 つの DOUBLE 要素は、それぞれ複素数の実部と虚部を表します。

      入力テンソルの DataType が DT_INT32、DT_UINT8、DT_INT16、DT_INT8、DT_QINT8、DT_QUINT8、DT_QINT32、DT_QINT16、DT_QUINT16、または DT_UINT16 の場合、content の要素型は INT です。

      入力テンソルの DataType が DT_INT64 の場合、content の要素型は LONG です。

      入力テンソルの DataType が DT_STRING の場合、content の要素型は STRING です。

      入力テンソルの DataType が DT_BOOL の場合、content の要素型は BOOLEAN です。

TFResponse クラス

TensorFlow モデルの予測結果から出力データを解析してアクセスします。

API

説明

List<Long> getTensorShape(String outputName)

  • 指定した出力テンソルの形状を取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:テンソルの形状を表す long のリストです。

List<Float> getFloatVals(String outputName)

  • 機能:出力テンソルのデータ型が DT_FLOAT、DT_COMPLEX64、DT_BFLOAT16、または DT_HALF の場合、指定した出力テンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

List<Double> getDoubleVals(String outputName)

  • 機能:出力テンソルのデータ型が DT_DOUBLE または DT_COMPLEX128 の場合、指定した出力テンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

List<Integer> getIntVals(String outputName)

  • 機能:出力テンソルのデータ型が DT_INT32、DT_UINT8、DT_INT16、DT_INT8、DT_QINT8、DT_QUINT8、DT_QINT32、DT_QINT16、DT_QUINT16、または DT_UINT16 の場合、指定した出力テンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

List<String> getStringVals(String outputName)

  • 機能:出力テンソルのデータ型が DT_STRING の場合、そのテンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

List<Long> getInt64Vals(String outputName)

  • 機能:出力テンソルのデータ型が DT_INT64 の場合、指定した出力テンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

List<Boolean> getBoolVals(String outputName)

  • 機能:出力テンソルのデータ型が DT_BOOL の場合、指定した出力テンソルのデータを取得します。

  • パラメーターoutputName は出力テンソルのエイリアスです。

  • 戻り値:モデルから出力された TensorData をフラット化した 1 次元配列です。

QueueClient class

Interacts with the EAS queue service to produce, consume, and manage data.

API

Description

QueueClient(String endpoint, String queueName, String token, HttpConfig httpConfig, QueueUser user)

  • Constructs a QueueClient instance.

  • Parameters:

    • endpoint: The endpoint address of the queue service.

    • queueName: The name of the queue service.

    • token: The token for service access.

    • httpConfig: The HTTP request configuration.

    • user: The user configuration. Specifies UserId (a random UUID by default) and GroupName (eas by default).

JSONObject attributes()

  • Gets detailed attributes of the queue service.

  • Returns: A JSONObject with queue information, including the following fields:

    • meta.maxPayloadBytes: The maximum allowed size (in bytes) for a single data item.

    • meta.name: The queue name.

    • stream.approxMaxLength: The approximate maximum number of items the queue can store.

    • stream.firstEntry: The index of the first item in the queue.

    • stream.lastEntry: The index of the last item in the queue.

    • stream.length: The current number of items in the queue.

Pair<Long, String> put(byte[] data, long priority, Map<String, String> tags)

  • Writes a data item to the queue.

  • Parameters:

    • data: The data to write, as a byte array.

    • priority: The data priority. 1 for high priority, 0 for normal priority (default).

    • tags: A map of custom key-value tags.

  • Returns: A Pair<Long, String> containing the index of the new data item and the request ID.

DataFrame[] get(long index, long length, long timeout, boolean autoDelete, Map<String, String> tags)

  • Retrieves data items from the queue.

  • Parameters:

    • index: The starting index from which to retrieve data. Use -1 to read the latest data.

    • length: The number of data items to retrieve.

    • timeout: The timeout period in seconds.

    • autoDelete: If true, the data is automatically deleted from the queue after being retrieved.

    • tags: A map of custom key-value tags, such as a RequestID.

  • Returns: An array of DataFrame objects.

void truncate(Long index)

  • Deletes all data items in the queue with an index less than the specified index.

String delete(Long index)

  • Deletes a specific data item from the queue.

  • Parameter: index is the index of the data item to delete.

  • Returns: "OK" on successful deletion.

JSONObject search(long index)

  • Queries the status of a specific data item in the queue.

  • Parameter: index is the index of the data item to query.

  • Returns: A JSONObject with queuing information, including:

    • ConsumerId: The ID of the instance processing the item.

    • IsPending: true if the item is being processed; false if it is waiting in the queue.

      • True means it is being processed.

      • False means it is queued.

    • WaitCount: The number of items ahead in the queue. This is valid only if IsPending is false. If IsPending is true, this value is 0.

    Example Responses:

    • The service returns {'ConsumerId': 'eas.****', 'IsPending': False, 'WaitCount':2}, which indicates that the request is being queued.

    • The log shows no data in stream and returns {}. This indicates that the data was not found in the queue. This may be because the data has been successfully processed by the server-side and a result has been returned, or the index parameter is incorrectly configured. Please check and confirm.

重要

When calling search, you must set the group ID in the QueueUser object to the service name. Otherwise, IsPending in the search result is always false.

  • Set group ID to the service name:

    QueueUser u = new QueueUser(UUID.randomUUID().toString(), "<service_name>");
    QueueClient input_queue = new QueueClient(queueEndpoint, inputQueueName, queueToken, new HttpConfig(), u);
  • Query the status for the specified index:

    System.out.println(input_queue.search(index));

WebSocketWatcher watch(long index, long window, boolean indexOnly, boolean autoCommit, Map<String, String> tags)

  • Subscribes to the queue service to receive data items as they become available.

  • Parameters:

    • index: The starting index. Use -1 to ignore all pending data and start with the newest items.

    • window: The size of the sending window (the maximum number of uncommitted items). The service will pause sending if the number of uncommitted items reaches this window size.

    • indexOnly: If true, returned DataFrame objects contain only the index and tags, not the data payload, to save bandwidth.

    • autoCommit: If true, items are automatically committed upon receipt, and the commit() call is not needed. When autoCommit is set to true, the window parameter is ignored.

    • tags: A map of custom parameters for the subscription request.

  • Returns: A WebSocketWatcher object for receiving the subscribed data. See the queue service example for usage details.

String commit(Long index) orString commit(Long[] index)

  • Confirms that one or more data items have been consumed, which deletes them from the queue.

  • Returns: "OK" on successful commit.

void end(boolean force)

Closes the connection to the queue service.

DataFrame クラス

キューサービスのデータ項目のラッパーです。

API

説明

byte[] getData()

  • データペイロードを取得します。

  • 戻り値:バイト配列のデータです。

long getIndex()

  • データ項目のインデックスを取得します。

  • 戻り値long 型のデータインデックスです。

Map<String, String> getTags()

  • データ項目に関連付けられているタグを取得します。

  • 戻り値: Tags という名前の Map<String,String> オブジェクト。たとえば、df.getTags().get("requestId") のようにして RequestID を取得できます。

コード例

同期推論の例

サービスの入出力形式に一致する例を選択してください。

String

カスタムプロセッサを使用してサービスをデプロイした場合、通常は文字列を使用して呼び出します。次の例に示すように、この方法は PMML モデルサービスで一般的に用いられます。

import com.aliyun.openservices.eas.predict.http.PredictClient;
import com.aliyun.openservices.eas.predict.http.HttpConfig;

public class TestString {
    public static void main(String[] args) throws Exception {
        // クライアントを初期化します。クライアントオブジェクトは共有してください。リクエストごとに新しいクライアントオブジェクトを作成しないでください。
        PredictClient client = new PredictClient(new HttpConfig());
        client.setToken("YWFlMDYyZDNmNTc3M2I3MzMwYmY0MmYwM2Y2MTYxMTY4NzBkNzdj****");
        // ダイレクトネットワーク接続を使用するには、setDirectEndpoint メソッドを呼び出します。
        // 例:client.setDirectEndpoint("182848887922****.vpc.cn-shanghai.aliyuncs.com");
        // ダイレクトネットワーク接続を有効にするには、EAS コンソールで有効化し、EAS サービスへのアクセスに使用するソース vSwitch を指定する必要があります。これによりゲートウェイをバイパスし、ソフトウェアロードバランシングを介してサービスインスタンスに直接アクセスできるため、安定性とパフォーマンスが向上します。
        // 注:標準のゲートウェイアクセスでは、ユーザー ID で始まるエンドポイントを使用します。このエンドポイントは、EAS コンソールのサービスの [Call Information] で確認できます。ダイレクトネットワーク接続では、182848887922****.vpc.{region_id}.aliyuncs.com の形式のドメイン名を使用します。
        client.setEndpoint("182848887922****.vpc.cn-shanghai.pai-eas.aliyuncs.com");
        client.setModelName("scorecard_pmml_example");

        // 入力文字列を定義します。
        String request = "[{\"money_credit\": 3000000}, {\"money_credit\": 10000}]";
        System.out.println(request);

        // EAS からレスポンス文字列を取得します。
        try {
            String response = client.predict(request);
            System.out.println(response);
        } catch (Exception e) {
            e.printStackTrace();
        }

        // クライアントをシャットダウンします。
        client.shutdown();
        return;
    }
}

TensorFlow

TensorFlow モデルを使用する場合は、次の例に示すように、入出力に TFRequest クラスと TFResponse クラスを使用します。

import java.util.List;

import com.aliyun.openservices.eas.predict.http.PredictClient;
import com.aliyun.openservices.eas.predict.http.HttpConfig;
import com.aliyun.openservices.eas.predict.request.TFDataType;
import com.aliyun.openservices.eas.predict.request.TFRequest;
import com.aliyun.openservices.eas.predict.response.TFResponse;

public class TestTF {
    public static TFRequest buildPredictRequest() {
        TFRequest request = new TFRequest();
        request.setSignatureName("predict_images");
        float[] content = new float[784];
        for (int i = 0; i < content.length; i++) {
            content[i] = (float) 0.0;
        }
        request.addFeed("images", TFDataType.DT_FLOAT, new long[]{1, 784}, content);
        request.addFetch("scores");
        return request;
    }

    public static void main(String[] args) throws Exception {
        PredictClient client = new PredictClient(new HttpConfig());

        // ダイレクトネットワーク接続を使用するには、setDirectEndpoint メソッドを呼び出します。エンドポイントの形式は {uid}.vpc.{region_id}.aliyuncs.com です。
        // client.setDirectEndpoint("182848887922****.vpc.cn-shanghai.aliyuncs.com");
        // 標準のゲートウェイアクセスでは、ユーザー ID で始まるエンドポイントを使用します。このエンドポイントは、EAS コンソールのサービスの [Call Information] で確認できます。
        client.setEndpoint("182848887922****.vpc.cn-shanghai.pai-eas.aliyuncs.com");
        client.setModelName("mnist_saved_model_example");
        client.setToken("YTg2ZjE0ZjM4ZmE3OTc0NzYxZDMyNmYzMTJjZTQ1YmU0N2FjMTAy****");
        long startTime = System.currentTimeMillis();
        int count = 1000;
        for (int i = 0; i < count; i++) {
            try {
                TFResponse response = client.predict(buildPredictRequest());
                List<Float> result = response.getFloatVals("scores");
                System.out.print("Predict Result: [");
                for (int j = 0; j < result.size(); j++) {
                    System.out.print(result.get(j).floatValue());
                    if (j != result.size() - 1) {
                        System.out.print(", ");
                    }
                }
                System.out.print("]\n");
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
        long endTime = System.currentTimeMillis();
        System.out.println("Spend Time: " + (endTime - startTime) + "ms");
        client.shutdown();
    }
}

キューサービスの例

キューサービスにアクセスするには、QueueClient インターフェイスを使用します。その方法を次の例に示します。

import com.alibaba.fastjson.JSONObject;
import com.aliyun.openservices.eas.predict.http.HttpConfig;
import com.aliyun.openservices.eas.predict.http.QueueClient;
import com.aliyun.openservices.eas.predict.queue_client.QueueUser;
import com.aliyun.openservices.eas.predict.queue_client.WebSocketWatcher;

public class DemoWatch {
    public static void main(String[] args) throws Exception {
        /** キューサービス クライアントを作成します。 */
        String queueEndpoint = "18*******.cn-hangzhou.pai-eas.aliyuncs.com";
        String inputQueueName = "test_queue_service";
        String sinkQueueName = "test_queue_service/sink";
        String queueToken = "test-token";

        /** 入力キュー。推論サービスはこのキューからリクエストデータを自動的に読み取ります。 */
        QueueClient inputQueue =
            new QueueClient(queueEndpoint, inputQueueName, queueToken, new HttpConfig(), new QueueUser());
        /** 出力キュー。推論サービスが入力データを処理した後、このキューに結果を書き込みます。 */
        QueueClient sinkQueue =
            new QueueClient(queueEndpoint, sinkQueueName, queueToken, new HttpConfig(), new QueueUser());
        /** キューデータをクリアします。取り扱いにご注意ください。 */
        inputQueue.clear();
        sinkQueue.clear();

        /** 入力キューにデータを追加します。 */
        int count = 10;
        for (int i = 0; i < count; ++i) {
            String data = Integer.toString(i);
            inputQueue.put(data.getBytes(), null);
            /** キューサービスは複数の優先度をサポートしています。put メソッドを使用してデータの優先度を設定できます。デフォルトの優先度は 0 です。 */
            //  inputQueue.put(data.getBytes(), 0L, null);
        }

        /** watch メソッドを使用して、出力キューのデータをサブスクライブします。ウィンドウサイズは 5 です。 */
        WebSocketWatcher watcher = sinkQueue.watch(0L, 5L, false, true, null);
        /** WatchConfig パラメーターを使用して、再試行回数、再試行間隔 (秒単位)、および無期限に再試行するかどうかをカスタマイズできます。WatchConfig を設定しない場合、システムはデフォルトで 5 秒間隔で 3 回再試行します。 */
        //  WebSocketWatcher watcher = sinkQueue.watch(0L, 5L, false, true, null, new WatchConfig(3, 1));
        //  WebSocketWatcher watcher = sinkQueue.watch(0L, 5L, false, true, null, new WatchConfig(true, 10));

        /** 出力データを取得します。 */
        for (int i = 0; i < count; ++i) {
            try {
                /** getDataFrame() メソッドは DataFrame データを取得します。この呼び出しは、データが利用可能になるまでブロックされます。 */
                byte[] data = watcher.getDataFrame().getData();
                System.out.println("[watch] data = " + new String(data));
            } catch (RuntimeException ex) {
                System.out.println("[watch] error = " + ex.getMessage());
                break;
            }
        }
        /** watcher オブジェクトを閉じます。各クライアントインスタンスは 1 つの watcher オブジェクトのみをサポートします。watcher を閉じない場合、次回の実行時にエラーが発生します。 */
        watcher.close();

        Thread.sleep(2000);
        JSONObject attrs = sinkQueue.attributes();
        System.out.println(attrs.toString());

        /** クライアントをシャットダウンします。 */
        inputQueue.shutdown();
        sinkQueue.shutdown();
    }
}

Java SDK を使用してサービスを呼び出す手順は次のとおりです:

  1. QueueClient インターフェイスを使用して、キューサービスのクライアントオブジェクトを作成します。キューサービスを使用する推論サービスでは、入力キューと出力キューのオブジェクトも作成する必要があります。

  2. put() 関数を使用して入力キューにデータを送信し、watch() 関数を使用して出力キューのデータをサブスクライブします。

    説明

    本番環境では、データの送信とデータのサブスクライブに別々のスレッドを使用してください。説明のため、この例では同一スレッドでこれらの操作を実行しています。

リクエストデータの圧縮

大量のデータを含むリクエストの場合、EAS は、データを Zlib または Gzip 形式で圧縮してからサーバーに送信することに対応しています。この機能を有効にするには、サービス構成で rpc.decompressor を指定する必要があります。

サービス構成は次のとおりです:

"metadata": {
  "rpc": {
    "decompressor": "zlib"
  }
}

コード例を次に示します:

package com.aliyun.openservices.eas.predict;
import com.aliyun.openservices.eas.predict.http.Compressor;
import com.aliyun.openservices.eas.predict.http.PredictClient;
import com.aliyun.openservices.eas.predict.http.HttpConfig;
public class TestString {
    public static void main(String[] args) throws Exception{
    	  // クライアントを初期化します。
        PredictClient client = new PredictClient(new HttpConfig());
        client.setEndpoint("18*******.cn-hangzhou.pai-eas.aliyuncs.com");
        client.setModelName("echo_compress");
        client.setToken("YzZjZjQwN2E4NGRkMDMxNDk5NzhhZDcwZDBjOTZjOGYwZDYxZGM2****");
        // Compressor.Gzip も使用できます。
        client.setCompressor(Compressor.Zlib);
        // 入力文字列を定義します。
        String request = "[{\"money_credit\": 3000000}, {\"money_credit\": 10000}]";
        System.out.println(request);
        // EAS からレスポンス文字列を取得します。
        String response = client.predict(request);
        System.out.println(response);
        // クライアントをシャットダウンします。
        client.shutdown();
        return;
    }
}

トラブルシューティング

Java SDK呼び出し例外 (認証ルーティング接続サーバー側のエラーなど) をトラブルシューティングするには、「Service Invocation SDK」の「呼び出し例外のトラブルシューティング」セクションをご参照ください。

サービスステータスコードエラーメッセージ の意味、および推奨されるアクションの一覧については、「付録:サービスステータスコードと一般的なエラー」をご参照ください。