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.
RemarqueNous vous recommandons d'utiliser localement une version de Python identique à celle préinstallée avec le moteur VVR cible.
RemarquePar 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 :
Sur la page , cliquez sur le nom du job cible.
-
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.0RemarqueAucun 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.17sélectionné lors du déploiement, installezapache-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.
Pour le développement de code métier Apache Flink V1.20, reportez-vous au Guide de développement Flink Python API.
Pour résoudre les problèmes rencontrés lors du codage Apache Flink, consultez la FAQ.
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 :
Connectez-vous à la console Realtime Compute.
Cliquez sur Console dans la colonne Actions de l'espace de travail cible.
Dans le volet de navigation de gauche, cliquez sur Artifacts.
-
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.
Sur la page , cliquez sur . Sélectionnez le package Python du connecteur cible dans le champ Additional Dependencies, configurez les autres paramètres, puis déployez le job.
-
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' -
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 :
Connectez-vous à la console Realtime Compute et accédez à l'espace de travail cible.
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.
-
Sur la page , cliquez sur 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.
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.
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
Pour un exemple complet de workflow de développement de jobs Flink Python, consultez Jobs Flink Python.
Pour utiliser des environnements virtuels Python personnalisés, des packages Python tiers, des fichiers JAR et des fichiers de données dans les jobs Flink Python, reportez-vous à Utiliser des dépendances Python.
Realtime Compute for Apache Flink prend également en charge les jobs SQL et DataStream. Pour plus d'informations, consultez le Guide de développement de jobs et Développer des jobs JAR.