Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Référence du SDK Python

Dernière mise à jour :Aug 09, 2026

Cette rubrique décrit comment installer et utiliser le SDK Python pour Realtime Compute for Apache Flink.

Prérequis

  • Créez une paire de clés AccessKey. Pour plus d'informations, consultez la section Créer une paire de clés AccessKey.

    Remarque

    Pour éviter les risques de sécurité liés à l'exposition de la paire de clés AccessKey de votre compte Alibaba Cloud, nous vous recommandons d'utiliser celle d'un utilisateur RAM. Créez un utilisateur RAM, accordez-lui les autorisations nécessaires pour accéder à Realtime Compute for Apache Flink, puis utilisez sa paire de clés AccessKey pour appeler le SDK. Pour plus d'informations, consultez les rubriques suivantes :

  • Python 3.6 ou une version ultérieure est installé.

  • Votre compte dispose des autorisations requises. Pour plus d'informations, consultez la section Gérer les autorisations.

Installer le SDK Python Flink

Installez le SDK Python à l'aide de pip.

  • Lors du développement de jobs et des opérations de maintenance (O&M), vous devez appeler l'API de la console de développement Realtime Compute. Pour plus de détails sur l'installation et l'utilisation, consultez le Centre SDK de la console de développement.

    pip3 install alibabacloud_ververica20220718==1.2.1
  • Pour afficher les informations sur les workspaces, en acheter ou ajuster les ressources, vous devez appeler l'API de la console de vente Realtime Compute. Pour plus d'informations sur l'installation et l'utilisation, consultez le Centre SDK de la console de vente Realtime Compute.

    pip3 install alibabacloud_foasconsole20211028==1.0.2

Tester les API et générer des exemples de SDK en ligne

OpenAPI Explorer simplifie l'utilisation des API. Utilisez-le pour effectuer des appels d'API, générer dynamiquement du code d'exemple pour le SDK et rechercher rapidement des opérations d'API afin d'optimiser votre processus de développement. Vous pouvez consulter et télécharger le code d'exemple du SDK sur les pages de référence des API de la console de développement et de la console de vente Realtime Compute. Pour connaître les étapes détaillées, consultez la section Démarrages rapides.

Dans la section Exemples de SDK, sélectionnez Python, puis cliquez sur Download Complete Project pour obtenir le projet d'exemple complet du SDK pour l'API.

Exemples de code

Remarque
  • Les endpoints de la console de vente Realtime Compute sont répertoriés dans la section Endpoints.

  • Les endpoints de la console de développement sont répertoriés dans la section Endpoints.

Afficher les workspaces achetés

Cet exemple montre comment interroger les détails des workspaces achetés dans une région spécifiée. Le paramètre de requête suivant est requis .

Region : ID de la région. Par exemple, 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:
        """
        Use an AccessKey pair to initialize the client.
        @return: Client
        @throws Exception
        """
        # Hard-coding your AccessKey pair into your project code can lead to security risks. We recommend using a more secure method, such as STS. The following code is for reference only.
        config = open_api_models.Config(
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set in your runtime environment.
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set in your runtime environment.
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # Modify the endpoint based on your actual requirements.
        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:
            # Call the API and print the response.
            response=client.describe_instances_with_options(describe_instances_request, runtime)
            print(response)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting 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:
            # If you copy this code to run, print the API response yourself.
            await client.describe_instances_with_options_async(describe_instances_request, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

Lister les déploiements

Cet exemple montre comment lister tous les déploiements dans un namespace. Les paramètres de requête suivants sont requis .

  • workspace : ID du workspace. Vous pouvez obtenir cet ID à partir du paramètre ResourceId renvoyé par l'exemple Afficher les workspaces achetés. Exemple : adf9e5147a****.

  • namespace : Nom du namespace. Exemple : 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:
        """
        Use an AccessKey pair to initialize the client.
        @return: Client
        @throws Exception
        """
        # Hard-coding your AccessKey pair into your project code can lead to security risks. We recommend using a more secure method, such as STS. The following code is for reference only.
        config = open_api_models.Config(
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set in your runtime environment.
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set in your runtime environment.
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # Modify the endpoint based on your actual requirements.
        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:
            # Call the API and print the response.
            request=client.list_deployments_with_options('namespace', list_deployments_request, list_deployments_headers, runtime)
            print(request)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting 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:
            # If you copy this code to run, print the API response yourself. The `namespace` parameter specifies the name of the namespace.
            await client.list_deployments_with_options_async('namespace', list_deployments_request, list_deployments_headers, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

Démarrer un job

Cet exemple montre comment démarrer un job à partir d'un déploiement dans un namespace. Les paramètres de requête suivants sont requis .

  • workspace : ID du workspace. Exemple : adf9e5147a****.

  • namespace : Nom du namespace. Exemple : script****-default.

  • deploymentId : ID du déploiement. Vous pouvez obtenir cet ID en appelant l'opération ListDeployments. Exemple : 3171d4d1-5952-4d02-b978-e762493b****.

  • kind : Type du décalage de démarrage. Valeurs valides : NONE (démarrage sans état), LATEST_SAVEPOINT (démarrage à partir du dernier point de sauvegarde), FROM_SAVEPOINT (démarrage à partir d'un point de sauvegarde spécifié) et LATEST_STATE (démarrage à partir du dernier état).

# -*- 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:
        """
        Use an AccessKey pair to initialize the client.
        @return: Client
        @throws Exception
        """
        # Hard-coding your AccessKey pair into your project code can lead to security risks. We recommend using a more secure method, such as STS. The following code is for reference only.
        config = open_api_models.Config(
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set in your runtime environment.
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set in your runtime environment.
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # Modify the endpoint based on your actual requirements.
        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:
            # If you copy this code to run, print the API response yourself.
            client.start_job_with_params_with_options('namespace', start_job_with_params_request, start_job_with_params_headers, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting 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(
            # The restore strategy for the job.
            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:
            # If you copy this code to run, print the API response yourself.
            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:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

Lister les jobs

Cet exemple montre comment lister tous les jobs pour un déploiement spécifique. Les paramètres de requête suivants sont requis .

  • workspace : ID du workspace. Exemple : adf9e5147a****.

  • namespace : Nom du namespace. Exemple : script****-default.

  • deploymentId : ID du déploiement. Vous pouvez obtenir cet ID en appelant l'opération ListDeployments. Exemple : 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:
        """
        Use an AccessKey pair to initialize the client.
        @return: Client
        @throws Exception
        """
        # Hard-coding your AccessKey pair into your project code can lead to security risks. We recommend using a more secure method, such as STS. The following code is for reference only.
        config = open_api_models.Config(
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set in your runtime environment.
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set in your runtime environment.
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # Modify the endpoint based on your actual requirements.
        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:
            # Call the API and print the response.
            request=client.list_jobs_with_options('namespace', list_jobs_request, list_jobs_headers, runtime)
            print(request)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting 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:
            # If you copy this code to run, print the API response yourself.
            await client.list_jobs_with_options_async('namespace', list_jobs_request, list_jobs_headers, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

Arrêter un job

Cet exemple montre comment arrêter un job. Les paramètres de requête suivants sont requis .

  • workspace : ID du workspace. Exemple : adf9e5147a****.

  • namespace : Nom du namespace. Exemple : script****-default.

  • jobId : ID du job. Vous pouvez obtenir cet ID en appelant l'opération ListJobs. Exemple : 3171d4d1-5952-4d02-b978-e762493b****.

  • stopStrategy : Stratégie d'arrêt. Valeurs valides : NONE (arrêt immédiat), STOP_WITH_SAVEPOINT (création d'un point de sauvegarde avant l'arrêt) et STOP_WITH_DRAIN (arrêt avec vidage).

# -*- 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:
        """
        Use an AccessKey pair to initialize the client.
        @return: Client
        @throws Exception
        """
        # Hard-coding your AccessKey pair into your project code can lead to security risks. We recommend using a more secure method, such as STS. The following code is for reference only.
        config = open_api_models.Config(
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_ID environment variable is set in your runtime environment.
            access_key_id=os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
            # Required. Ensure that the ALIBABA_CLOUD_ACCESS_KEY_SECRET environment variable is set in your runtime environment.
            access_key_secret=os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
        )
        # Modify the endpoint based on your actual requirements.
        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(
            # The stop strategy for the job.
            stop_strategy='stopStrategy'
        )
        stop_job_request = ververica_20220718_models.StopJobRequest(
            body=stop_job_request_body
        )
        runtime = util_models.RuntimeOptions()
        try:
            # If you copy this code to run, print the API response yourself.
            client.stop_job_with_options('namespace', 'jobId', stop_job_request, stop_job_headers, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting 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:
            # If you copy this code to run, print the API response yourself.
            await client.stop_job_with_options_async('namespace', 'jobId', stop_job_request, stop_job_headers, runtime)
        except Exception as error:
            # This is for demonstration only. Implement proper error handling in your production code and do not ignore exceptions.
            # Error message
            print(error.message)
            # Troubleshooting URL
            print(error.data.get("Recommend"))
            UtilClient.assert_as_string(error.message)
if __name__ == '__main__':
    Sample.main(sys.argv[1:])

Documentation connexe

Pour plus de détails sur le SDK Java, consultez la section Référence du SDK Java.