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

Realtime Compute for Apache Flink:Flink Java SDK

最終更新日:Jun 22, 2026

このトピックでは、Realtime Compute for Apache Flink Java SDK をインストールして使用する方法を説明します。

注意事項

2022年 9月 19日から 10月 27日にかけて、Alibaba Cloud は全リージョンで Realtime Compute for Apache Flink の SDK のアップデートをリリースしました。この新しいバージョンがデフォルトになりました。

説明
  • SDK アップデートの影響については、「お知らせ」をご参照ください。

  • このドキュメントでは、新しい SDK バージョンについて説明します。以前のバージョンのドキュメントを参照するには、「OpenAPI SDK (非推奨)」をクリックしてください。

前提条件

  • AccessKey が作成されていること。詳細については、「AccessKey の作成」をご参照ください。

    説明

    Alibaba Cloud アカウント の AccessKey が漏洩した場合のセキュリティリスクを軽減するために、RAM ユーザー を作成し、Flink 関連の権限を付与してから、その AccessKey を使用して SDK を呼び出すことを推奨します。詳細については、以下のトピックをご参照ください。

  • Java 8 以降がインストールされていること。

  • アカウントに必要な権限があること。詳細については、「権限管理」をご参照ください。

Realtime Compute for Apache Flink Java SDK

商用コンソール

インストール方法

コード

Apache Maven

<dependency>

<groupId>com.aliyun</groupId>

<artifactId>foasconsole20211028</artifactId>

<version>2.1.0</version>

</dependency>

Gradle Groovy DSL

implementation 'com.aliyun:foasconsole20211028:2.1.0'

Gradle Kotlin DSL

implementation("com.aliyun:foasconsole20211028:2.1.0")

Scala SBT

libraryDependencies += "com.aliyun" % "foasconsole20211028" % "2.1.0"

Apache Ivy

<dependency org="com.aliyun" name="foasconsole20211028" rev="2.1.0" />

Groovy Grape

@Grapes(

@Grab(group='com.aliyun', module='foasconsole20211028', version='2.1.0')

)

Leiningen

[com.aliyun/foasconsole20211028 "2.1.0"]

Apache Buildr

'com.aliyun:foasconsole20211028:jar:2.1.0'

開発コンソール

インストール方法

コード

Apache Maven

<dependency>

<groupId>com.aliyun</groupId>

<artifactId>ververica20220718</artifactId>

<version>1.7.0</version>

</dependency>

Gradle Groovy DSL

implementation 'com.aliyun:ververica20220718:1.7.0'

Gradle Kotlin DSL

implementation("com.aliyun:ververica20220718:1.7.0")

Scala SBT

libraryDependencies += "com.aliyun" % "ververica20220718" % "1.7.0"

Apache Ivy

<dependency org="com.aliyun" name="ververica20220718" rev="1.7.0" />

Groovy Grape

@Grapes(

@Grab(group='com.aliyun', module='ververica20220718', version='1.7.0')

)

Leiningen

[com.aliyun/ververica20220718 "1.7.0"]

Apache Buildr

'com.aliyun:ververica20220718:jar:1.7.0'

オンラインデバッグと SDK サンプルの生成

OpenAPI Explorer を使用して、API オペレーションをオンラインで呼び出し、SDK サンプルコードを動的に生成し、API オペレーションを迅速に検索できます。Realtime Compute for Apache Flink ページと Realtime Compute Selling Console ページで、API オペレーションの SDK サンプルコードを表示してダウンロードできます。詳細については、「クイックスタート」をご参照ください。

[SDK サンプル] タブで、[Java] などの言語を選択し、[サンプルの実行] をクリックしてコードをオンラインで実行するか、[完全なプロジェクトのダウンロード] をクリックしてサンプルコードをダウンロードします。

サンプル

説明
  • Realtime Compute for Apache Flink 商用コンソールのエンドポイントについては、エンドポイントをご参照ください。

  • Realtime Compute for Apache Flink 開発コンソールのエンドポイントについては、エンドポイントをご参照ください。

購入済みワークスペース

このサンプルでは、指定したリージョン内の購入済み Realtime Compute for Apache Flink ワークスペースの詳細を取得する方法を説明します。

Region:リージョンの ID。詳細については、「エンドポイント」をご参照ください。例:cn-hangzhou。

package com.aliyun.sample;
import com.aliyun.foasconsole20211028.models.DescribeInstancesResponse;
import com.aliyun.tea.*;
import com.alibaba.fastjson2.JSON;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.foasconsole20211028.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "foasconsole.aliyuncs.com";
        return new com.aliyun.foasconsole20211028.Client(config);
    }
    public static void main(String[] args_) throws Exception {
        com.aliyun.foasconsole20211028.Client client = Sample.createClient();
        com.aliyun.foasconsole20211028.models.DescribeInstancesRequest describeInstancesRequest = new com.aliyun.foasconsole20211028.models.DescribeInstancesRequest()
                .setRegion("cn-beijing");
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        try {
            DescribeInstancesResponse response = client.describeInstancesWithOptions(describeInstancesRequest, runtime);
            System.out.println(response.statusCode);
            // インスタンスのゾーン ID を取得します。
            System.out.println(response.getBody().getInstances().get(0).zoneId);
            // インスタンスが属するリソースグループの ID を取得します。
            System.out.println(response.getBody().getInstances().get(0).resourceGroupId);
            System.out.println(JSON.toJSON(response));
        } catch (TeaException error) {
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        } catch (Exception _error) {
            TeaException error = new TeaException(_error.getMessage(), _error);
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        }
    }
}

デプロイメントの作成

SQL デプロイメント

このサンプルでは、SQL デプロイメントを作成する方法を説明します。

  • workspace:ワークスペースの ID。この ID は、購入済みワークスペースの表示 時に返される ResourceId から取得できます。例:adf9e514****

  • namespace:名前空間の名前。例:test-default。

  • body.name:ジョブの名前。例:mysql_data_holo_test。

  • body.engineVersion:エンジンバージョン。例:vvr-8.0.7-flink-1.17。サポートされているエンジンバージョンの一覧表示 により、サポートされているエンジンバージョンを取得できます。

  • body.sqlArtifact.sqlScript:SQL スクリプトの内容。例:CREATE TEMPORARY TABLE datagen_source( name VARCHAR ) WITH ( 'connector' = 'datagen' ); CREATE TEMPORARY TABLE blackhole_sink( name VARCHAR ) with ( 'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT name from datagen_source;

  • body.sqlArtifact.kind:ジョブの種類。例:SQLSCRIPT。

  • body.deploymentTarget.mode:デプロイモード。PER_JOB モードのみがサポートされています。

  • body.deploymentTarget.name:デプロイキューの名前。例:default-queue。

  • body.executionMode:実行モード。例:STREAMING (ストリーミングモード)。

  • body.streamingResourceSetting.resourceSettingMode:ストリーミングモードのリソースモード。例:BASIC。

  • body.streamingResourceSetting.basicResourceSetting.jobmanagerResourceSettingSpec.cpu:JobManager の CPU コア数。例:2。

  • body.streamingResourceSetting.basicResourceSetting.jobmanagerResourceSettingSpec.memory:JobManager のメモリ。例:4 GiB。

  • body.streamingResourceSetting.basicResourceSetting.taskmanagerResourceSettingSpec.cpu:TaskManager の CPU コア数。例:2。

  • body.streamingResourceSetting.basicResourceSetting.taskmanagerResourceSettingSpec.memory:TaskManager のメモリ。例:4 GiB。

package com.aliyun.sample;
import com.aliyun.tea.*;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.teaopenapi.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-beijing.aliyuncs.com";
        return new com.aliyun.teaopenapi.Client(config);
    }
    /**
     * <b>description</b> :
     * <p>API パラメータ</p>
     *
     * @param path params
     * @return OpenApi.Params
     */
    public static com.aliyun.teaopenapi.models.Params createApiInfo(String namespace) throws Exception {
        com.aliyun.teaopenapi.models.Params params = new com.aliyun.teaopenapi.models.Params()
                // API オペレーションの名前。
                .setAction("CreateDeployment")
                // API オペレーションのバージョン。
                .setVersion("2022-07-18")
                // API オペレーションのプロトコル。
                .setProtocol("HTTPS")
                // API オペレーションの HTTP メソッド。
                .setMethod("POST")
                .setAuthType("AK")
                .setStyle("ROA")
                // API オペレーションのリクエストパス。
                .setPathname("/api/v2/namespaces/" + namespace + "/deployments")
                // リクエストボディの形式。
                .setReqBodyType("json")
                // レスポンスボディの形式。
                .setBodyType("json");
        return params;
    }
    public static void main(String[] args_) throws Exception {
        java.util.List<String> args = java.util.Arrays.asList(args_);
        com.aliyun.teaopenapi.Client client = Sample.createClient();
        com.aliyun.teaopenapi.models.Params params = Sample.createApiInfo("test-default");
        // body params
        java.util.Map<String, Object> body = TeaConverter.buildMap(
                new TeaPair("name", "mysql_data_holo_test"),
                new TeaPair("engineVersion", "vvr-8.0.7-flink-1.17"),
                new TeaPair("artifact", TeaConverter.buildMap(
                        new TeaPair("sqlArtifact", TeaConverter.buildMap(
                                new TeaPair("sqlScript", "CREATE TEMPORARY TABLE datagen_source(   name VARCHAR ) WITH (   'connector' = 'datagen' ); CREATE TEMPORARY TABLE blackhole_sink(   name  VARCHAR ) with (   'connector' = 'blackhole' ); INSERT INTO blackhole_sink SELECT name from datagen_source;")
                        )),
                        new TeaPair("kind", "SQLSCRIPT")
                )),
                new TeaPair("deploymentTarget", TeaConverter.buildMap(
                        new TeaPair("mode", "PER_JOB"),
                        new TeaPair("name", "default-queue")
                )),
                new TeaPair("executionMode", "STREAMING"),
                new TeaPair("streamingResourceSetting", TeaConverter.buildMap(
                        new TeaPair("resourceSettingMode", "BASIC"),
                        new TeaPair("basicResourceSetting", TeaConverter.buildMap(
                                new TeaPair("jobmanagerResourceSettingSpec", TeaConverter.buildMap(
                                        new TeaPair("cpu", 2),
                                        new TeaPair("memory", "4Gi")
                                )),
                                new TeaPair("taskmanagerResourceSettingSpec", TeaConverter.buildMap(
                                        new TeaPair("cpu", 2),
                                        new TeaPair("memory", "4Gi")
                                ))
                        ))
                ))
        );
        // header params
        java.util.Map<String, String> headers = new java.util.HashMap<>();
        headers.put("workspace", "ab2*******884d");
        // runtime options
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        com.aliyun.teaopenapi.models.OpenApiRequest request = new com.aliyun.teaopenapi.models.OpenApiRequest()
                .setHeaders(headers)
                .setBody(body);
        // このメソッドは Map を返します。マップからレスポンスボディ、レスポンスヘッダー、および HTTP ステータスコードを取得できます。
        client.callApi(params, request, runtime);
        java.util.Map<String, ?> response = client.callApi(params, request, runtime);
        System.out.println(response);
    }
}

JAR デプロイメント

このサンプルでは、JAR デプロイメントを作成する方法を説明します。

説明
  • JAR パッケージを OSS バケットにアップロードし、Realtime Compute for Apache Flink ワークスペースからアクセスできることを確認してください。詳細については、「シンプルアップロード」をご参照ください。

  • ファイルのアップロード後、ダウンロード URL は https://<Bucket>.oss-<Region>.aliyuncs.com/<FileName> の形式になります。

  • workspace:ワークスペースの ID。この ID は、購入済みワークスペースの表示 時に返される ResourceId から取得できます。例:adf9e514****。

  • namespace:名前空間の名前。例:test-default。

  • body.name:ジョブ名。例:my-test-jar。

  • body.engineVersion:エンジンバージョン。例:vvr-8.0.7-flink-1.17。サポートされているエンジンバージョンの一覧を取得する ことができます。

  • body.jarArtifact.kind:ジョブアーティファクトの種類。例:JAR。

  • body.jarArtifact.jarUri:JAR デプロイメントの完全な URL パス。例:https://myBucket.oss-cn-hangzhou.aliyuncs.com/test.jar。

  • body.jarArtifact.entryClass:エントリポイントクラス。完全修飾クラス名を指定する必要があります。例:org.apache.flink.test。

  • body.deploymentTarget.mode:デプロイモード。PER_JOB モードのみがサポートされています。

  • body.deploymentTarget.name:デプロイキューの名前。例:default-queue。

  • body.executionMode:実行モード。例:STREAMING (ストリーミングモード)。

  • body.streamingResourceSetting.resourceSettingMode:ストリーミングモードのリソースモード。例:BASIC。

  • body.streamingResourceSetting.basicResourceSetting.jobmanagerResourceSettingSpec.cpu:JobManager の CPU コア数。例:2。

  • body.streamingResourceSetting.basicResourceSetting.jobmanagerResourceSettingSpec.memory:JobManager のメモリ。例:4 GiB。

  • body.streamingResourceSetting.basicResourceSetting.taskmanagerResourceSettingSpec.cpu:TaskManager の CPU コア数。例:2。

  • body.streamingResourceSetting.basicResourceSetting.taskmanagerResourceSettingSpec.memory:TaskManager のメモリ。例:4 GiB。

package com.aliyun.sample;
import com.aliyun.tea.*;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.teaopenapi.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-hangzhou.aliyuncs.com";
        return new com.aliyun.teaopenapi.Client(config);
    }
    /**
     * <b>description</b> :
     * <p>API パラメータ</p>
     *
     * @param path params
     * @return OpenApi.Params
     */
    public static com.aliyun.teaopenapi.models.Params createApiInfo(String namespace) throws Exception {
        com.aliyun.teaopenapi.models.Params params = new com.aliyun.teaopenapi.models.Params()
                // API オペレーションの名前。
                .setAction("CreateDeployment")
                // API オペレーションのバージョン。
                .setVersion("2022-07-18")
                // API オペレーションのプロトコル。
                .setProtocol("HTTPS")
                // API オペレーションの HTTP メソッド。
                .setMethod("POST")
                .setAuthType("AK")
                .setStyle("ROA")
                // API オペレーションのリクエストパス。
                .setPathname("/api/v2/namespaces/" + namespace + "/deployments")
                // リクエストボディの形式。
                .setReqBodyType("json")
                // レスポンスボディの形式。
                .setBodyType("json");
        return params;
    }
    public static void main(String[] args_) throws Exception {
        java.util.List<String> args = java.util.Arrays.asList(args_);
        com.aliyun.teaopenapi.Client client = Sample.createClient();
        com.aliyun.teaopenapi.models.Params params = Sample.createApiInfo("flink-default");
        // body params
        java.util.Map<String, Object> body = TeaConverter.buildMap(
                new TeaPair("name", "my-test-jar"),
                new TeaPair("engineVersion", "vvr-8.0.7-flink-1.17"),
                new TeaPair("artifact", TeaConverter.buildMap(
                        new TeaPair("kind", "JAR"),
                        new TeaPair("jarArtifact", TeaConverter.buildMap(
                                new TeaPair("jarUri", "https://flink-test.oss-cn-hangzhou.aliyuncs.com/flinkDemo.jar?*****"),
                                new TeaPair("entryClass", "com.aliyun.FlinkDemo")
                        ))
                )),
                new TeaPair("deploymentTarget", TeaConverter.buildMap(
                        new TeaPair("mode", "PER_JOB"),
                        new TeaPair("name", "default-queue")
                )),
                new TeaPair("executionMode", "STREAMING"),
                new TeaPair("streamingResourceSetting", TeaConverter.buildMap(
                        new TeaPair("resourceSettingMode", "BASIC"),
                        new TeaPair("basicResourceSetting", TeaConverter.buildMap(
                                new TeaPair("jobmanagerResourceSettingSpec", TeaConverter.buildMap(
                                        new TeaPair("cpu", 2),
                                        new TeaPair("memory", "4Gi")
                                )),
                                new TeaPair("taskmanagerResourceSettingSpec", TeaConverter.buildMap(
                                        new TeaPair("cpu", 2),
                                        new TeaPair("memory", "4Gi")
                                ))
                        ))
                ))
        );
        // header params
        java.util.Map<String, String> headers = new java.util.HashMap<>();
        headers.put("workspace", "d05a*****e44");
        // runtime options
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        com.aliyun.teaopenapi.models.OpenApiRequest request = new com.aliyun.teaopenapi.models.OpenApiRequest()
                .setHeaders(headers)
                .setBody(body);
        // このメソッドは Map を返します。マップからレスポンスボディ、レスポンスヘッダー、および HTTP ステータスコードを取得できます。
        java.util.Map<String, ?> response = client.callApi(params, request, runtime);
        System.out.println(response);
    }
}

デプロイメントの一覧表示

このサンプルでは、名前空間内のすべてのデプロイメントを一覧表示する方法を説明します。

  • workspace:ワークスペースの ID。この ID は、購入済みワークスペースの表示 オペレーションで返される ResourceId から取得できます。例:adf9e514****。

  • namespace:名前空間の名前。例:test-default。

package com.aliyun.sample;
import com.aliyun.tea.*;
import com.alibaba.fastjson2.JSON;
import com.aliyun.ververica20220718.models.ListDeploymentsResponse;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.ververica20220718.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-hangzhou.aliyuncs.com";
        return new com.aliyun.ververica20220718.Client(config);
    }
    public static void main(String[] args_) throws Exception {
        com.aliyun.ververica20220718.Client client = Sample.createClient();
        com.aliyun.ververica20220718.models.ListDeploymentsHeaders listDeploymentsHeaders = new com.aliyun.ververica20220718.models.ListDeploymentsHeaders()
                .setWorkspace("ab2a******884d");
        com.aliyun.ververica20220718.models.ListDeploymentsRequest listDeploymentsRequest = new com.aliyun.ververica20220718.models.ListDeploymentsRequest();
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        try {
            ListDeploymentsResponse response=client.listDeploymentsWithOptions("test-default", listDeploymentsRequest, listDeploymentsHeaders, runtime);
            System.out.println(response.body.data.get(0).name);
            System.out.println(response.body.data.get(0).deploymentId);
            System.out.println(JSON.toJSON(response));
        } catch (TeaException error) {
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        } catch (Exception _error) {
            TeaException error = new TeaException(_error.getMessage(), _error);
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        }
    }
}

ジョブの開始

このサンプルでは、名前空間内のデプロイメントからジョブを開始する方法を説明します。

  • workspace:ワークスペース ID。例:adf9e5147a****。

  • namespace:名前空間の名前。例:test-default。

  • deploymentId:デプロイメントの ID。この ID は、デプロイメントの一覧表示 オペレーションを呼び出すことで取得できます。例:10283a02-****-****-****-8dabf617d52f。

  • kind:復元ストラテジーの種類。サポートされている値は、NONE (ステートレス開始)、LATEST_SAVEPOINT (最新のセーブポイントから開始)、FROM_SAVEPOINT (指定されたセーブポイントから開始)、LATEST_STATE (最新の状態から開始) です。

package com.aliyun.sample;
import com.aliyun.tea.*;
import com.aliyun.ververica20220718.models.StartJobWithParamsResponse;
import com.alibaba.fastjson2.JSON;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.ververica20220718.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-hangzhou.aliyuncs.com";
        return new com.aliyun.ververica20220718.Client(config);
    }
    public static void main(String[] args_) throws Exception {
        com.aliyun.ververica20220718.Client client = Sample.createClient();
        com.aliyun.ververica20220718.models.StartJobWithParamsHeaders startJobWithParamsHeaders = new com.aliyun.ververica20220718.models.StartJobWithParamsHeaders()
                .setWorkspace("ab2a******884d");
        com.aliyun.ververica20220718.models.DeploymentRestoreStrategy jobStartParametersDeploymentRestoreStrategy = new com.aliyun.ververica20220718.models.DeploymentRestoreStrategy()
                .setKind("NONE");
        com.aliyun.ververica20220718.models.JobStartParameters jobStartParameters = new com.aliyun.ververica20220718.models.JobStartParameters()
                .setRestoreStrategy(jobStartParametersDeploymentRestoreStrategy)
                .setDeploymentId("10283a02-****-****-****-8dabf617d52f");
        com.aliyun.ververica20220718.models.StartJobWithParamsRequest startJobWithParamsRequest = new com.aliyun.ververica20220718.models.StartJobWithParamsRequest()
                .setBody(jobStartParameters);
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        try {
            StartJobWithParamsResponse response = client.startJobWithParamsWithOptions("test-default", startJobWithParamsRequest, startJobWithParamsHeaders, runtime);
            System.out.println(JSON.toJSON(response.body));
        } catch (TeaException error) {
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        } catch (Exception _error) {
            TeaException error = new TeaException(_error.getMessage(), _error);
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        }
    }
}

ジョブの一覧表示

このサンプルでは、特定のデプロイメントのすべてのジョブを取得する方法を説明します。

  • workspace:ワークスペース ID。例:adf9e5147****。

  • namespace:名前空間の名前。例:test-default。

  • deploymentId:デプロイメント ID。この ID を取得するには、「デプロイメントの一覧表示」をご参照ください。例:8489b7ec-****-****-****-cc4c17fa12b0。

package com.aliyun.sample;
import com.aliyun.tea.*;
import com.aliyun.ververica20220718.models.ListJobsResponse;
import com.alibaba.fastjson2.JSON;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.ververica20220718.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-beijing.aliyuncs.com";
        return new com.aliyun.ververica20220718.Client(config);
    }
    public static void main(String[] args_) throws Exception {
        com.aliyun.ververica20220718.Client client = Sample.createClient();
        com.aliyun.ververica20220718.models.ListJobsHeaders listJobsHeaders = new com.aliyun.ververica20220718.models.ListJobsHeaders()
                .setWorkspace("ab2a******884d");
        com.aliyun.ververica20220718.models.ListJobsRequest listJobsRequest = new com.aliyun.ververica20220718.models.ListJobsRequest()
                .setDeploymentId("8489b7ec-****-****-****-cc4c17fa12b0");
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        try {
            ListJobsResponse response =  client.listJobsWithOptions("test-default", listJobsRequest, listJobsHeaders, runtime);
            // ジョブの実行結果を表示します。
            System.out.println("Execution result is: "+response.body.success);
            // ジョブ ID を取得します。この ID はジョブの停止に使用できます。
            System.out.println(response.body.getData().get(0).jobId);
            System.out.println(JSON.toJSON(response));
        } catch (TeaException error) {
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        } catch (Exception _error) {
            TeaException error = new TeaException(_error.getMessage(), _error);
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        }
    }
}

ジョブの停止

このサンプルでは、ジョブを停止する方法を説明します。

  • workspace:ワークスペース ID。例:adf9e5147****。

  • namespace:名前空間の名前。例:test-default。

  • jobId:ジョブの ID。この ID は、ジョブの一覧表示 オペレーションを呼び出すことで取得できます。例:3171d4d1-****-****-****-e762493b7765。

  • stopStrategy:ジョブ停止ストラテジー。サポートされている値は、NONE (ジョブを直接停止)、STOP_WITH_SAVEPOINT (セーブポイントの作成後にジョブを停止)、STOP_WITH_DRAIN (ドレインモードでジョブを停止) です。

package com.aliyun.sample;
import com.alibaba.fastjson2.JSON;
import com.aliyun.tea.*;
import com.aliyun.ververica20220718.models.StopJobResponse;
public class Sample {
    /**
     * <b>description</b> :
     * <p>AccessKey ID と AccessKey Secret を使用してクライアントを初期化します。</p>
     * @return Client
     *
     * @throws Exception
     */
    public static com.aliyun.ververica20220718.Client createClient() throws Exception {
        // コード内に AccessKey ペアをハードコーディングすると、AccessKey ペアが漏洩し、アカウント内のすべてのリソースのセキュリティが脅かされる可能性があります。以下のサンプルコードは参考用です。
        com.aliyun.teaopenapi.models.Config config = new com.aliyun.teaopenapi.models.Config()
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_ID 環境変数が設定されていることを確認してください。
                .setAccessKeyId(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"))
                // 必須。ランタイム環境で ALIBABA_CLOUD_ACCESS_KEY_SECRET 環境変数が設定されていることを確認してください。
                .setAccessKeySecret(System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        config.endpoint = "ververica.cn-hangzhou.aliyuncs.com";
        return new com.aliyun.ververica20220718.Client(config);
    }
    public static void main(String[] args_) throws Exception {
        java.util.List<String> args = java.util.Arrays.asList(args_);
        com.aliyun.ververica20220718.Client client = Sample.createClient();
        com.aliyun.ververica20220718.models.StopJobHeaders stopJobHeaders = new com.aliyun.ververica20220718.models.StopJobHeaders()
                .setWorkspace("ab2a******884d");
        com.aliyun.ververica20220718.models.StopJobRequestBody stopJobRequestBody = new com.aliyun.ververica20220718.models.StopJobRequestBody()
                .setStopStrategy("NONE");
        com.aliyun.ververica20220718.models.StopJobRequest stopJobRequest = new com.aliyun.ververica20220718.models.StopJobRequest()
                .setBody(stopJobRequestBody);
        com.aliyun.teautil.models.RuntimeOptions runtime = new com.aliyun.teautil.models.RuntimeOptions();
        try {
            StopJobResponse response = client.stopJobWithOptions("test-default", "7970e881-****-****-****-1a3746710878", stopJobRequest, stopJobHeaders, runtime);
            System.out.println(JSON.toJSON(response.getBody().getData()));
        } catch (TeaException error) {
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        } catch (Exception _error) {
            TeaException error = new TeaException(_error.getMessage(), _error);
            // このサンプルでは、エラーメッセージを参考として出力します。実際のプロジェクトでは、例外を無視せず、慎重に処理することを推奨します。
            // エラーメッセージ
            System.out.println(error.getMessage());
            // 診断アドレス
            System.out.println(error.getData().get("Recommend"));
            com.aliyun.teautil.Common.assertAsString(error.message);
        }
    }
}

参考資料

Python SDK の詳細については、「Python SDK リファレンス」をご参照ください。