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

Realtime Compute for Apache Flink:Python SDK リファレンス

最終更新日:Jun 22, 2026

このトピックでは、Realtime Compute for Apache Flink の Python SDK のインストール方法と使用方法について説明します。

前提条件

  • AccessKey ペアを作成します。 詳細については、AccessKey ペアの作成をご参照ください。

    説明

    Alibaba Cloud アカウント (root ユーザー) の AccessKey ペアが漏洩するとセキュリティリスクが生じるため、RAM ユーザーの AccessKey ペアを使用することを推奨します。RAM ユーザーを作成し、Realtime Compute for Apache Flink にアクセスするために必要な権限をユーザーに付与してから、そのユーザーの AccessKey ペアを使用して SDK を呼び出します。 詳細については、次のトピックをご参照ください。

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

  • ご利用のアカウントに必要な権限が付与されていること。 詳細については、権限の管理をご参照ください。

Flink Python SDK のインストール

pip を使用して Python SDK をインストールします。

  • ジョブ開発や O&M などの操作を実行する場合、Realtime Compute 開発コンソール API を呼び出す必要があります。 インストールと使用方法の詳細については、開発コンソール SDK センターをご参照ください。

    pip3 install alibabacloud_ververica20220718==1.2.1
  • ワークスペース情報の表示、ワークスペースの購入、またはリソースの調整を行うには、Realtime Compute 販売コンソール API を呼び出す必要があります。 インストールと使用方法の詳細については、Realtime Compute 販売コンソール SDK センターをご参照ください。

    pip3 install alibabacloud_foasconsole20211028==1.0.2

API のテストと SDK サンプルコードのオンライン生成

OpenAPI Explorer を使用すると、API の使用が簡単になります。 API 呼び出しの実行、SDK サンプルコードの動的な生成、API オペレーションの迅速な検索が可能になり、開発プロセスを効率化できます。 開発コンソールおよびRealtime Compute 販売コンソールの API リファレンスページで SDK サンプルコードを表示およびダウンロードできます。 詳細な手順については、クイックスタートをご参照ください。

SDK サンプルセクションで Python を選択し、[Download Complete Project] をクリックすると、API の完全な SDK サンプルプロジェクトを取得できます。

コード例

説明
  • Realtime Compute 販売コンソールのエンドポイントは、エンドポイントに記載されています。

  • 開発コンソールのエンドポイントは、エンドポイントに記載されています。

購入済みワークスペースの表示

この例では、指定されたリージョンで購入したワークスペースの詳細をクエリする方法を示します。 次のリクエストパラメーターは必須ですをご参照ください。

Region:リージョン ID。 例:cn-hangzhou。をご参照ください。

# -*- coding: utf-8 -*-
import os
import sys
from typing import List
from alibabacloud_foasconsole20211028.client import Client as foasconsole20211028Client
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_foasconsole20211028 import models as foasconsole_20211028_models
from alibabacloud_tea_util import models as util_models
from alibabacloud_tea_util.client import Client as UtilClient
class Sample:
    def __init__(self):
        pass
    @staticmethod
    def create_client() -> foasconsole20211028Client:
        """
        AccessKey ペアを使用してクライアントを初期化します。
        @return: Client
        @throws Exception
        """
        # AccessKey ペアをプロジェクトコードにハードコーディングすると、セキュリティリスクにつながる可能性があります。 STS などのより安全な方法を使用することを推奨します。 以下のコードは参照用です。
        config = open_api_models.Config(
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID が設定されていることを確認してください。
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_SECRET が設定されていることを確認してください。
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # 実際にご利用の要件に基づいてエンドポイントを変更します。
        config.endpoint = f'foasconsole.aliyuncs.com'
        return foasconsole20211028Client(config)
    @staticmethod
    def main(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        describe_instances_request = foasconsole_20211028_models.DescribeInstancesRequest(
            region='cn-hangzhou'
        )
        runtime = util_models.RuntimeOptions()
        try:
            # API を呼び出して応答を出力します。
            response=client.describe_instances_with_options(describe_instances_request, runtime)
            print(response)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
    @staticmethod
    async def main_async(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        describe_instances_request = foasconsole_20211028_models.DescribeInstancesRequest(
            region='cn-hangzhou'
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            await client.describe_instances_with_options_async(describe_instances_request, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

デプロイメントのリスト表示

この例では、名前空間内のすべてのデプロイメントをリスト表示する方法を示します。 次のリクエストパラメーターは必須ですをご参照ください。

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

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

# -*- coding: utf-8 -*-
import os
import sys
from typing import List
from alibabacloud_ververica20220718.client import Client as ververica20220718Client
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_ververica20220718 import models as ververica_20220718_models
from alibabacloud_tea_util import models as util_models
from alibabacloud_tea_util.client import Client as UtilClient
class Sample:
    def __init__(self):
        pass
    @staticmethod
    def create_client() -> ververica20220718Client:
        """
        AccessKey ペアを使用してクライアントを初期化します。
        @return: Client
        @throws Exception
        """
        # AccessKey ペアをプロジェクトコードにハードコーディングすると、セキュリティリスクにつながる可能性があります。 STS などのより安全な方法を使用することを推奨します。 以下のコードは参照用です。
        config = open_api_models.Config(
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID が設定されていることを確認してください。
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_SECRET が設定されていることを確認してください。
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # 実際にご利用の要件に基づいてエンドポイントを変更します。
        config.endpoint = f'ververica.cn-hangzhou.aliyuncs.com'
        return ververica20220718Client(config)
    @staticmethod
    def main(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        list_deployments_headers = ververica_20220718_models.ListDeploymentsHeaders(
            workspace='workspace'
        )
        list_deployments_request = ververica_20220718_models.ListDeploymentsRequest()
        runtime = util_models.RuntimeOptions()
        try:
            # API を呼び出して応答を出力します。
            request=client.list_deployments_with_options('namespace', list_deployments_request, list_deployments_headers, runtime)
            print(request)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
    @staticmethod
    async def main_async(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        list_deployments_headers = ververica_20220718_models.ListDeploymentsHeaders(
            workspace='workspace'
        )
        list_deployments_request = ververica_20220718_models.ListDeploymentsRequest()
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。 `namespace` パラメーターは名前空間の名前を指定します。
            await client.list_deployments_with_options_async('namespace', list_deployments_request, list_deployments_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

ジョブの開始

この例では、名前空間のデプロイメントからジョブを開始する方法を示します。 次のリクエストパラメーターは必須ですをご参照ください。

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

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

  • deploymentId:デプロイメント ID。 この ID は、ListDeployments 操作を呼び出すことで取得できます。 例:3171d4d1-5952-4d02-b978-e762493b****。

  • kind:開始オフセットのタイプ。 有効な値:NONE (ステートレス開始)、LATEST_SAVEPOINT (最新のセーブポイントから開始)、FROM_SAVEPOINT (指定されたセーブポイントから開始)、および LATEST_STATE (最新の状態から開始)。

# -*- coding: utf-8 -*-
import os
import sys
from typing import List
from alibabacloud_ververica20220718.client import Client as ververica20220718Client
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_ververica20220718 import models as ververica_20220718_models
from alibabacloud_tea_util import models as util_models
from alibabacloud_tea_util.client import Client as UtilClient
class Sample:
    def __init__(self):
        pass
    @staticmethod
    def create_client() -> ververica20220718Client:
        """
        AccessKey ペアを使用してクライアントを初期化します。
        @return: Client
        @throws Exception
        """
        # AccessKey ペアをプロジェクトコードにハードコーディングすると、セキュリティリスクにつながる可能性があります。 STS などのより安全な方法を使用することを推奨します。 以下のコードは参照用です。
        config = open_api_models.Config(
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID が設定されていることを確認してください。
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_SECRET が設定されていることを確認してください。
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # 実際にご利用の要件に基づいてエンドポイントを変更します。
        config.endpoint = f'ververica.cn-hangzhou.aliyuncs.com'
        return ververica20220718Client(config)
    @staticmethod
    def main(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        start_job_with_params_headers = ververica_20220718_models.StartJobWithParamsHeaders(
            workspace='workspace'
        )
        job_start_parameters_deployment_restore_strategy = ververica_20220718_models.DeploymentRestoreStrategy(
            kind='NONE'
        )
        job_start_parameters = ververica_20220718_models.JobStartParameters(
            deployment_id='deploymentId',
            restore_strategy=job_start_parameters_deployment_restore_strategy
        )
        start_job_with_params_request = ververica_20220718_models.StartJobWithParamsRequest(
            body=job_start_parameters
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            client.start_job_with_params_with_options('namespace', start_job_with_params_request, start_job_with_params_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
    @staticmethod
    async def main_async(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        start_job_with_params_headers = ververica_20220718_models.StartJobWithParamsHeaders(
            workspace='workspace'
        )
        job_start_parameters_deployment_restore_strategy = ververica_20220718_models.DeploymentRestoreStrategy(
            # ジョブの復元戦略。
            kind='NONE'
        )
        job_start_parameters = ververica_20220718_models.JobStartParameters(
            deployment_id='deploymentId',
            restore_strategy=job_start_parameters_deployment_restore_strategy
        )
        start_job_with_params_request = ververica_20220718_models.StartJobWithParamsRequest(
            body=job_start_parameters
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            await client.start_job_with_params_with_options_async('namespace', start_job_with_params_request, start_job_with_params_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

ジョブのリスト表示

この例では、特定のデプロイメントのすべてのジョブをリスト表示する方法を示します。 次のリクエストパラメーターは必須ですをご参照ください。

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

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

  • deploymentId:デプロイメント ID。 この ID は、ListDeployments 操作を呼び出すことで取得できます。 例:3171d4d1-5952-4d02-b978-e762493b****。

# -*- coding: utf-8 -*-
import os
import sys
from typing import List
from alibabacloud_ververica20220718.client import Client as ververica20220718Client
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_ververica20220718 import models as ververica_20220718_models
from alibabacloud_tea_util import models as util_models
from alibabacloud_tea_util.client import Client as UtilClient
class Sample:
    def __init__(self):
        pass
    @staticmethod
    def create_client() -> ververica20220718Client:
        """
        AccessKey ペアを使用してクライアントを初期化します。
        @return: Client
        @throws Exception
        """
        # AccessKey ペアをプロジェクトコードにハードコーディングすると、セキュリティリスクにつながる可能性があります。 STS などのより安全な方法を使用することを推奨します。 以下のコードは参照用です。
        config = open_api_models.Config(
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID が設定されていることを確認してください。
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_SECRET が設定されていることを確認してください。
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # 実際にご利用の要件に基づいてエンドポイントを変更します。
        config.endpoint = f'ververica.cn-hangzhou.aliyuncs.com'
        return ververica20220718Client(config)
    @staticmethod
    def main(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        list_jobs_headers = ververica_20220718_models.ListJobsHeaders(
            workspace='workspace'
        )
        list_jobs_request = ververica_20220718_models.ListJobsRequest(
            deployment_id='deploymentId'
        )
        runtime = util_models.RuntimeOptions()
        try:
            # API を呼び出して応答を出力します。
            request=client.list_jobs_with_options('namespace', list_jobs_request, list_jobs_headers, runtime)
            print(request)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
    @staticmethod
    async def main_async(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        list_jobs_headers = ververica_20220718_models.ListJobsHeaders(
            workspace='workspace'
        )
        list_jobs_request = ververica_20220718_models.ListJobsRequest(
            deployment_id='deploymentId'
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            await client.list_jobs_with_options_async('namespace', list_jobs_request, list_jobs_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

ジョブの停止

この例では、ジョブを停止する方法を示します。 次のリクエストパラメーターは必須ですをご参照ください。

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

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

  • jobId:ジョブ ID。 この ID は、ListJobs 操作を呼び出すことで取得できます。 例:3171d4d1-5952-4d02-b978-e762493b****。

  • stopStrategy:停止戦略。 有効な値:NONE (即時停止)、STOP_WITH_SAVEPOINT (停止前にセーブポイントを作成)、および STOP_WITH_DRAIN (ドレインして停止)。

# -*- coding: utf-8 -*-
import os
import sys
from typing import List
from alibabacloud_ververica20220718.client import Client as ververica20220718Client
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_ververica20220718 import models as ververica_20220718_models
from alibabacloud_tea_util import models as util_models
from alibabacloud_tea_util.client import Client as UtilClient
class Sample:
    def __init__(self):
        pass
    @staticmethod
    def create_client() -> ververica20220718Client:
        """
        AccessKey ペアを使用してクライアントを初期化します。
        @return: Client
        @throws Exception
        """
        # AccessKey ペアをプロジェクトコードにハードコーディングすると、セキュリティリスクにつながる可能性があります。 STS などのより安全な方法を使用することを推奨します。 以下のコードは参照用です。
        config = open_api_models.Config(
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_ID が設定されていることを確認してください。
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # 必須。 実行環境に環境変数 ALIBABA_CLOUD_ACCESS_KEY_SECRET が設定されていることを確認してください。
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # 実際にご利用の要件に基づいてエンドポイントを変更します。
        config.endpoint = f'ververica.cn-hangzhou.aliyuncs.com'
        return ververica20220718Client(config)
    @staticmethod
    def main(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        stop_job_headers = ververica_20220718_models.StopJobHeaders(
            workspace='workspace'
        )
        stop_job_request_body = ververica_20220718_models.StopJobRequestBody(
            # ジョブの停止戦略。
            stop_strategy='stopStrategy'
        )
        stop_job_request = ververica_20220718_models.StopJobRequest(
            body=stop_job_request_body
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            client.stop_job_with_options('namespace', 'jobId', stop_job_request, stop_job_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
    @staticmethod
    async def main_async(
        args: List[str],
    ) -> None:
        client = Sample.create_client()
        stop_job_headers = ververica_20220718_models.StopJobHeaders(
            workspace='workspace'
        )
        stop_job_request_body = ververica_20220718_models.StopJobRequestBody(
            stop_strategy='stopStrategy'
        )
        stop_job_request = ververica_20220718_models.StopJobRequest(
            body=stop_job_request_body
        )
        runtime = util_models.RuntimeOptions()
        try:
            # このコードをコピーして実行する場合は、API の応答を自身で出力してください。
            await client.stop_job_with_options_async('namespace', 'jobId', stop_job_request, stop_job_headers, runtime)
        except Exception as error:
            # これはデモ用です。 本番コードでは適切なエラー処理を実装し、例外を無視しないでください。
            # エラーメッセージ
            print(error.message)
            # トラブルシューティング URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

関連ドキュメント

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