Tous les produits
Search
Centre de documentation

Realtime Compute for Apache Flink:Développer des jobs Python

Dernière mise à jour :Aug 12, 2026

Cette rubrique présente le contexte, les limites, les méthodes de développement et de débogage, ainsi que l'utilisation des connecteurs pour les jobs Flink Python API.

Contexte

Développez vos jobs Flink Python en local. Une fois le développement terminé, déployez et démarrez les jobs sur la console de développement Flink pour obtenir les résultats métier. Pour plus d'informations sur la procédure globale, consultez Jobs Flink Python.

Prérequis de l'environnement de développement

  • Les versions du moteur Realtime Compute for Apache Flink antérieures à VVR 8.0.11 intègrent Python 3.7.9. Les versions VVR 8.0.11 et ultérieures incluent Python 3.9.21. À partir de VVR 11.7, les versions Python 3.9.21, 3.10.19 et 3.11.14 sont préinstallées.

    Remarque

    Nous vous recommandons d'utiliser localement une version de Python identique à celle préinstallée avec le moteur VVR cible.

    Remarque

    Par défaut, VVR 11.7 et les versions ultérieures utilisent Python 3.9. Pour utiliser les environnements Python 3.10 ou 3.11 préinstallés, configurez le job comme suit :

    1. Sur la page O&M > Deployments, cliquez sur le nom du job cible.

    2. Sous l'onglet Configuration, dans la section Parameters, cliquez sur Edit à droite, puis ajoutez la configuration suivante dans le champ Other Configuration (exemple avec Python 3,11) :

      python.executable: python3.11
      python.client.executable: python3.11
  • Installez PyFlink dans une version correspondant au moteur VVR cible. Par exemple, si le moteur sélectionné sur la page de déploiement est vvr-11.7.0-jdk11-flink-1.20, installez la version appropriée :

    pip install ververica-flink==11.7.0
    Remarque

    Aucun package PyFlink dédié n'est fourni pour VVR 11.5 et les versions antérieures. Installez la version correspondante de PyFlink open source. Ainsi, pour un moteur vvr-8.0.9-flink-1.17 sélectionné lors du déploiement, installez apache-flink==1.17.*.

    pip install apache-flink==1.17.2
  • Un IDE est nécessaire. Nous recommandons PyCharm ou VS Code.

  • Développez vos jobs Python en local avant de les déployer et de les exécuter via la console Realtime Compute for Apache Flink.

Limites

En raison des contraintes liées aux environnements de déploiement et réseau de Flink, respectez les limitations suivantes lors du développement de jobs Python :

  • Seules les versions open source Flink V1.13 et ultérieures sont prises en charge.

  • L'espace de travail Flink dispose d'un environnement Python préinstallé comprenant des bibliothèques courantes telles que Pandas, NumPy et PyArrow. Pour plus de détails, reportez-vous à la Liste des packages préinstallés à la fin de cette rubrique.

  • L'environnement d'exécution Flink prend uniquement en charge JDK 8 et JDK 11. Si votre job Python dépend de packages JAR tiers, assurez-vous de leur compatibilité.

  • VVR 4.x ne prend en charge que Scala V2.11 open source, tandis que VVR 6.x et les versions ultérieures requièrent Scala V2.12 open source. Vérifiez que les dépendances des packages JAR tiers correspondent à la version Scala appropriée.

Développement de jobs

Choix entre Table API/SQL et DataStream API

PyFlink propose deux approches de développement : Table API/SQL et DataStream API. Table API/SQL est recommandé pour les raisons suivantes :

  • Meilleures performances : Le plan d'exécution optimisé de Table API/SQL s'exécute entièrement au sein de la JVM. À l'inverse, DataStream API nécessite une sérialisation et une désérialisation ligne par ligne entre la JVM et le processus Python, ce qui engendre une surcharge significative.

  • Fonctionnalités plus complètes : Table API/SQL offre une prise en charge étendue des connecteurs, formats de données et fonctions de fenêtrage, tout en partageant l'écosystème de connecteurs des jobs SQL.

  • Recommandation communautaire : La communauté Apache Flink privilégie Table API/SQL comme axe principal de développement pour PyFlink.

Réservez DataStream API aux scénarios nécessitant une logique personnalisée complexe impossible à exprimer en SQL.

Références de développement

Consultez les documents suivants pour développer votre code métier Flink en local. Une fois le développement terminé, importez le code sur la console de développement Flink et déployez le job.

Structure du projet

La structure de projet recommandée pour les jobs Python est la suivante :

my-flink-python-project/
├── my_job.py                # Main job file.
├── udfs.py                  # User-defined functions (optional)
├── requirements.txt         # Third-party Python dependencies (optional)
└── config.properties        # Configuration file (optional)

Gestion des dépendances

Pour savoir comment utiliser des environnements virtuels Python personnalisés, des packages Python tiers, des fichiers JAR et des fichiers de données dans vos jobs Python, consultez Utiliser des dépendances Python.

Fonctions définies par l'utilisateur (UDF)

Voici un exemple de développement d'une UDSF Python permettant de masquer des données sensibles sous forme de chaîne :

from pyflink.table import DataTypes
from pyflink.table.udf import udf

@udf(result_type=DataTypes.STRING())def mask_phone(phone: str):"""Phone number masking: retains the first 3 and last 4 digits, replacing the middle with ****"""if phone is None or len(phone) != 11:return phone.
    return phone[:3] + '****' + phone[7:]

Utilisez cette UDF dans un job SQL :

CREATE TEMPORARY FUNCTION mask_phone AS 'udfs.mask_phone' LANGUAGE PYTHON;INSERT INTO sink_table
SELECT name, mask_phone(phone) AS masked_phone
FROM source_table;

Pour enregistrer, mettre à jour ou supprimer des UDF, reportez-vous à Gérer les UDF.

Utilisation des connecteurs

Pour obtenir la liste des connecteurs pris en charge par Flink, consultez Connecteurs pris en charge. Pour utiliser un connecteur, procédez comme suit :

  1. Connectez-vous à la console Realtime Compute.

  2. Cliquez sur Console dans la colonne Actions de l'espace de travail cible.

  3. Dans le volet de navigation de gauche, cliquez sur Artifacts.

  4. Cliquez sur Upload Artifact et sélectionnez le package Python du connecteur souhaité pour le charger.

    Vous pouvez charger un connecteur développé en interne ou fourni par Flink. Pour télécharger les packages Python officiels des connecteurs Flink, consultez la liste des connecteurs.

  5. Sur la page O&M > Deployments, cliquez sur Create Deployment > Python Deployment. Sélectionnez le package Python du connecteur cible dans le champ Additional Dependencies, configurez les autres paramètres, puis déployez le job.

  6. Cliquez sur le nom du job déployé. Sous l'onglet Configuration, dans la section Parameters, cliquez sur Edit. Dans Other Configuration, ajoutez les informations d'emplacement des packages Python du connecteur.

    Si votre job dépend de plusieurs packages Python de connecteurs, par exemple deux packages nommés connector-1.jar et connector-2.jar, utilisez la configuration suivante.

    pipeline.classpaths: 'file:///flink/usrlib/connector-1.jar;file:///flink/usrlib/connector-2.jar'
  7. Pour utiliser des connecteurs intégrés, des formats de données et des catalogues (uniquement pour VVR 11.2 et ultérieur), ajoutez la configuration dans la section Running Parameters sous Additional Configuration du job. Exemple :

    ## Multiple connectors used.
    pipeline.used-builtin-connectors: kafka;sls
    ## Multiple data formats for data transmission.
    pipeline.used-builtin-formats: avro;parquet
    ## Multiple previously created catalogs used.
    pipeline.used-builtin-catalogs: catalogname1;catalogname2

Pour plus de détails sur l'utilisation des connecteurs, consultez Exemple de code complet.

Débogage des jobs

Intégrez des instructions de journalisation dans le code des UDF Python afin de générer des journaux utiles au dépannage. Exemple :

import logging

@udf(result_type=DataTypes.BIGINT())
def add(i, j):
  logging.info("hello world")
  return i + j

Une fois les journaux générés, consultez-les dans les fichiers de journalisation du TaskManager.

Débogage local

Realtime Compute for Apache Flink n'ayant pas accès à Internet par défaut, votre code pourrait ne pas pouvoir se connecter directement aux sources de données en ligne pour des tests locaux. Nous recommandons les approches suivantes pour le débogage local :

  • Tests unitaires : Réalisez des tests unitaires indépendants sur les UDF afin de valider la logique des fonctions.

  • Exécution locale : Simulez les entrées à l'aide de sources de données locales (fichiers ou données en mémoire) et exécutez le job localement pour vérifier la logique de traitement. Exemple :

    from pyflink.datastream import StreamExecutionEnvironment
    
    env = StreamExecutionEnvironment.get_execution_environment()# Use a local data source for testing.
    ds = env.from_collection([('Alice', 1), ('Bob', 2), ('Alice', 3)])
    ds.key_by(lambda x: x[0]).sum(1).print()
    env.execute("local_test")
  • Débogage distant : Pour déboguer en vous connectant à des sources de données en ligne, consultez Exécuter et déboguer localement des jobs avec connecteurs.

Déploiement des jobs

Une fois le développement du job Python terminé, chargez-le sur la console Realtime Compute pour le déployer. Procédez comme suit :

  1. Connectez-vous à la console Realtime Compute et accédez à l'espace de travail cible.

  2. Dans le volet de navigation de gauche, cliquez sur Artifacts et chargez le fichier du job Python (.py ou .zip). Chargez également les dépendances tierces ou fichiers de configuration si nécessaire.

  3. Sur la page O&M > Deployments, cliquez sur Create Deployment > Python Deployment et renseignez les informations de déploiement.

    Parameter

    Description

    Python File Path

    Sélectionnez le fichier du job Python précédemment chargé.

    Entry Module

    Ce champ est facultatif pour un fichier .py. Pour un fichier .zip, spécifiez le nom du module d'entrée, par exemple my_job.

    Additional Dependencies

    Sélectionnez les fichiers JAR de connecteurs ou de configuration le cas échéant.

    Python Libraries

    Sélectionnez les packages Python tiers (.whl ou .zip) si nécessaire.

    Python Archives

    Sélectionnez un environnement virtuel Python personnalisé (.zip) le cas échéant.

  4. Cliquez sur Deploy.

Pour plus d'informations sur les paramètres de déploiement, consultez Déployer un job.

Exemple de code complet

Cet exemple illustre un job streaming Python qui lit des données depuis Kafka, effectue un traitement simple et écrit les résultats dans MySQL. Il est fourni à titre indicatif uniquement.

Remarque

Cet exemple n'inclut pas la configuration des paramètres d'exécution tels que les checkpoints et les stratégies de redémarrage. Vous pouvez personnaliser ces configurations sur la page Deployment Details après avoir déployé le job. Pour plus d'informations, consultez Configurer les informations de déploiement d'un job.

import logging
import sys

from pyflink.common import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment

logging.basicConfig(stream=sys.stdout, level=logging.INFO)

def kafka_to_mysql():
    # Create the execution environment.
    env = StreamExecutionEnvironment.get_execution_environment()
    t_env = StreamTableEnvironment.create(env)

    # Create the Kafka source table.
    t_env.execute_sql("""
        CREATE TABLE kafka_source (
            `id` INT,
            `name` STRING,
            `score` INT,
            `event_time` TIMESTAMP(3),
            WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
        ) WITH (
            'connector' = 'kafka',
            'topic' = 'student_topic',
            'properties.bootstrap.servers' = 'your-kafka-broker:9092',
            'properties.group.id' = 'my-group',
            'scan.startup.mode' = 'latest-offset',
            'format' = 'json'
        )
    """)

    # Create the MySQL sink table.
    t_env.execute_sql("""
        CREATE TABLE mysql_sink (
            `id` INT,
            `name` STRING,
            `score` INT,
            PRIMARY KEY (id) NOT ENFORCED
        ) WITH (
            'connector' = 'jdbc',
            'url' = 'jdbc:mysql://your-mysql-host:3306/my_database',
            'table-name' = 'student',
            'username' = 'your_username',
            'password' = 'your_password'
        )
    """)

    # Filter records with scores greater than 60 and write them to MySQL.
    t_env.execute_sql("""
        INSERT INTO mysql_sink
        SELECT id, name, score
        FROM kafka_source
        WHERE score >= 60
    """)

if __name__ == '__main__':
    kafka_to_mysql()

Packages préinstallés

VVR-11

Les packages suivants sont préinstallés dans l'espace de travail Flink.

Package

Version

apache-beam

2.48.0

avro-python3

1.10.2

brotlipy

0.7.0

certifi

2022.12.7

cffi

1.15.1

charset-normalizer

2.0.4

cloudpickle

2.2.1

conda

22.11.1

conda-content-trust

0.1.3

conda-package-handling

1.9.0

crcmod

1,7

cryptography

38.0.1

Cython

3.0.12

dill

0.3.1.1

dnspython

2.7.0

docopt

0.6.2

exceptiongroup

1.3.0

fastavro

1.12.1

fasteners

0,20

find_libpython

0.5.0

grpcio

1.56.2

grpcio-tools

1.56.2

hdfs

2.7.3

httplib2

0.22.0

idna

3,4

importlib_metadata

8.7.0

iniconfig

2.1.0

isort

6.1.0

numpy

1.24.4

objsize

0.6.1

orjson

3.9.15

packaging

25,0

pandas

2.3.3

pemja

0.5.5

pip

22.3.1

pluggy

1.0.0

proto-plus

1.26.1

protobuf

4.25.8

py-spy

0.4.0

py4j

0.10.9.7

pyarrow

11.0.0

pyarrow-hotfix

0,6

pycodestyle

2.14.0

pycosat

0.6.4

pycparser

2,21

pydot

1.4.2

pymongo

4.15.4

pyOpenSSL

22.0.0

pyparsing

3.2.5

PySocks

1.7.1

pytest

7.4.4

python-dateutil

2.9.0

pytz

2025.2

regex

2025.11.3

requests

2.32.5

ruamel.yaml

0.18.16

ruamel.yaml.clib

0.2.14

setuptools

70.0.0

six

1.16.0

tomli

2.3.0

toolz

0.12.0

tqdm

4.64.1

typing_extensions

4.15.0

tzdata

2025.2

urllib3

1.26.13

wheel

0.38.4

zipp

3.23.0

zstandard

0.25.0

torch

2.5.1

torchvision

0.20.1

transformers

4.57.6

opencv-python-headless

4.10.0.84

pillow

11.3.0

ultralytics

8.4.66

easyocr

1.7.2

open-clip-torch

2.32.0

rembg

2.0.61

onnxruntime

1.16.3

av

14.2.0

librosa

0.11.0

soundfile

0.13.1

imagehash

4.3.2

safetensors

0.7.0

huggingface-hub

0.36.2

tokenizers

0.22.2

timm

1.0.27

scikit-learn

1.6.1

scikit-image

0.24.0

scipy

1.13.1

numba

0.60.0

llvmlite

0.43.0

matplotlib

3.9.4

pywavelets

1.6.0

tifffile

2024.8.30

shapely

2.0.7

pymatting

1.1.15

audioread

3.1.0

soxr

0.5.0.post1

joblib

1.5.3

threadpoolctl

3.6.0

sympy

1.13.1

mpmath

1.3.0

python-bidi

0.6.10

filelock

3.19.1

pyyaml

6.0.3

decorator

5.3.1

msgpack

1.1.2

ftfy

6.3.1

wcwidth

0.8.1

contourpy

1.3.0

cycler

0.12.1

fonttools

4.60.2

importlib-resources

5.4.0

kiwisolver

1.4.7

lazy-loader

0,5

attrs

26.1.0

jsonschema

4.25.1

jsonschema-specifications

2025.9.1

referencing

0.36.2

rpds-py

0.27.1

pooch

1.9.0

platformdirs

4.4.0

VVR-8

Les packages suivants sont préinstallés dans l'espace de travail Flink.

Package

Version

apache-beam

2.43.0

avro-python3

1.9.2.1

certifi

2025.7.9

charset-normalizer

3.4.2

cloudpickle

2.2.0

crcmod

1,7

Cython

0.29.24

dill

0.3.1.1

docopt

0.6.2

fastavro

1.4.7

fasteners

0,19

find_libpython

0.4.1

grpcio

1.46.3

grpcio-tools

1.46.3

hdfs

2.7.3

httplib2

0.20.4

idna

3,10

isort

6.0.1

numpy

1.21.6

objsize

0.5.2

orjson

3.10.18

pandas

1.3.5

pemja

0.3.2

pip

22.3.1

proto-plus

1.26.1

protobuf

3.20.3

py4j

0.10.9.7

pyarrow

8.0.0

pycodestyle

2.14.0

pydot

1.4.2

pymongo

3.13.0

pyparsing

3.2.3

python-dateutil

2.9.0

pytz

2025.2

regex

2024.11.6

requests

2.32.4

setuptools

58.1.0

six

1.17.0

typing_extensions

4.14.1

urllib3

2.5.0

wheel

0.33.4

zstandard

0.23.0

VVR-6

Les packages suivants sont préinstallés dans l'espace de travail Flink.

Package

Version

apache-beam

2.27.0

avro-python3

1.9.2.1

certifi

2024.8.30

charset-normalizer

3.3.2

cloudpickle

1.2.2

crcmod

1,7

Cython

0.29.16

dill

0.3.1.1

docopt

0.6.2

fastavro

0.23.6

future

0.18.3

grpcio

1.29.0

hdfs

2.7.3

httplib2

0.17.4

idna

3,8

importlib-metadata

6.7.0

isort

5.11.5

jsonpickle

2.0.0

mock

2.0.0

numpy

1.19.5

oauth2client

4.1.3

pandas

1.1.5

pbr

6.1.0

pemja

0.1.4

pip

20.1.1

protobuf

3.17.3

py4j

0.10.9.3

pyarrow

2.0.0

pyasn1

0.5.1

pyasn1-modules

0.3.0

pycodestyle

2.10.0

pydot

1.4.2

pymongo

3.13.0

pyparsing

3.1.4

python-dateutil

2.8.0

pytz

2024.1

requests

2.31.0

rsa

4,9

setuptools

47.1.0

six

1.16.0

typing-extensions

3.7.4.3

urllib3

2.0.7

wheel

0.42.0

zipp

3.15.0

Références