Une fonction scalaire définie par l'utilisateur (UDSF) associe zéro, une ou plusieurs valeurs scalaires à une seule valeur scalaire. Chaque ligne en entrée génère exactement une valeur en sortie.
Cette rubrique explique comment créer, enregistrer et utiliser une UDSF Python dans Realtime Compute for Apache Flink.
Limites
Les contraintes suivantes s'appliquent au développement de fonctions définies par l'utilisateur (UDF) Python dans Realtime Compute for Apache Flink :
| Contrainte | Exigence |
|---|---|
| Version d'Apache Flink | 1.12 et versions ultérieures |
| Version de Python | Préinstallée sur chaque espace de travail. VVR antérieure à 8.0.11 : Python 3.7.9. VVR 8.0.11 et versions ultérieures : Python 3.9.21. |
| Version du JDK | JDK 8 et JDK 11. Les packages JAR tiers doivent être compatibles avec JDK 8 ou JDK 11. |
| Version de Scala | Uniquement Scala 2.11 open source. Les packages JAR tiers doivent être compatibles avec Scala 2.11. |
| Fonctions inline | Prise en charge uniquement dans VVR 11.9-preview1 et versions ultérieures. |
Après une mise à niveau vers VVR 8.0.11 ou version ultérieure, testez, déployez et exécutez à nouveau vos brouillons PyFlink existants afin de confirmer leur compatibilité.
Méthodes de développement
Vous disposez de deux méthodes pour développer une UDSF Python :
Empaquetez votre code Python, téléchargez le package, puis enregistrez-le en tant que fonction sur la plateforme.
Déclarez la logique du code sous forme de fonction inline dans une instruction SQL.
Si la logique de votre UDSF est simple, privilégiez le développement sous forme de fonction inline.
Créer une UDSF
Les étapes suivantes prennent Windows comme environnement d'exemple. Flink fournit un référentiel d'exemples qui inclut des implémentations pour les UDSF, les fonctions d'agrégation définies par l'utilisateur (UDAF) et les fonctions table définies par l'utilisateur (UDTF).
-
Téléchargez et décompressez python_demo-master sur votre machine locale.
Il s'agit d'un référentiel GitHub tiers. L'accès peut être lent ou intermittent.
Dans PyCharm, choisissez File > Open et ouvrez le répertoire
python_demo-masterdécompressé.-
Ouvrez
udfs.pydans le chemin\python_demo-master\udxet définissez votre UDSF.from pyflink.table import DataTypes from pyflink.table.udf import udf @udf(result_type=DataTypes.STRING()) def sub_string(s: str, begin: int, end: int): return s[begin:end]L'exemple
sub_stringextrait les caractères de la positionbeginà la positionenddans la chaîne d'entrée. -
Depuis le répertoire
\python_demo-master, exécutez la commande suivante pour empaqueter le répertoireudx:zip -r python_demo.zip udxLorsque
python_demo.zipapparaît dans\python_demo-master\, le package est prêt.
Enregistrer une UDSF
Une fois le package créé, enregistrez l'UDSF dans la console Realtime Compute for Apache Flink. Pour connaître les étapes d'enregistrement, consultez Gérer les fonctions définies par l'utilisateur (UDF).
Utiliser une UDSF
Après avoir enregistré l'UDSF, utilisez-la dans un job Flink SQL.
-
Créez un brouillon à l'aide de Flink SQL. Pour plus de détails, consultez Présentation du développement de jobs. L'exemple suivant appelle
ASI_UDSF(le nom enregistré de votre UDSF) pour extraire les caractères des positions 2 à 4 du champadans la table source :CREATE TEMPORARY TABLE ASI_UDSF_Source ( a VARCHAR, b INT, c INT ) WITH ( 'connector' = 'datagen' ); CREATE TEMPORARY TABLE ASI_UDSF_Sink ( a VARCHAR ) WITH ( 'connector' = 'blackhole' ); INSERT INTO ASI_UDSF_Sink SELECT ASI_UDSF(a, 2, 4) FROM ASI_UDSF_Source; Dans le volet de navigation de gauche de la console de développement, choisissez O&M > Deployments. Localisez le déploiement, puis cliquez sur Start dans la colonne Actions. Une fois le déploiement démarré, les caractères aux positions 2 à 4 du champ
adansASI_UDSF_Sourcesont écrits dansASI_UDSF_Sink.
Fonctions scalaires inline Python
Une fonction inline intègre son implémentation directement dans une instruction CREATE FUNCTION, ce qui permet de définir et d'enregistrer la fonction dans la même instruction SQL. L'exemple ci-dessous définit une fonction inline qui masque les adresses e-mail. Déclarez l'intégralité de la logique du code entre les délimiteurs $$.
CREATE TEMPORARY FUNCTION mask_email(email STRING)
RETURNS STRING AS $$
if email is None:
return None
name, separator, domain = email.partition("@")
if not separator:
return email
return name[:1] + "***@" + domain
$$ LANGUAGE PYTHON;
Python est sensible à l'indentation. Commencez le code de premier niveau du corps de la fonction à la colonne 0 et maintenez une indentation relative correcte au sein des blocs de code.
Les fonctions inline ne prennent pas en charge l'exécution vectorisée et ne peuvent pas être déclarées comme fonctions non déterministes.
Seules les fonctions TEMPORARY sont prises en charge. Les définitions de fonctions ne peuvent pas encore être persistées sur la plateforme.
Fonctions définies par l'utilisateur asynchrones
Pour les UDF effectuant des opérations intensives en E/S, telles que l'accès à des bases de données externes ou des requêtes HTTP, utilisez des fonctions définies par l'utilisateur asynchrones. Une seule UDF asynchrone peut traiter plusieurs requêtes d'E/S simultanément, répartissant ainsi le temps d'attente entre les requêtes et améliorant le débit du job.
Limites
-
Prise en charge uniquement sur VVR 11.7 et versions ultérieures. VVR PyFlink 11.7 ou version ultérieure est requis. Pour plus de détails, consultez ververica-flink.
pip3 install "ververica-flink>=11.7" Seules les fonctions scalaires définies par l'utilisateur (UDSF) asynchrones sont prises en charge.
Seul le mode processus Python est pris en charge, c'est-à-dire
python.execution-mode=process.Les fonctions définies par l'utilisateur asynchrones Pandas ne sont pas encore prises en charge.
Les fonctions inline ne sont pas encore prises en charge.
Utilisation
Une fonction définie par l'utilisateur asynchrone peut être implémentée sous forme de fonction async Python ou comme sous-classe de la classe de fonction asynchrone. Le code d'exemple est présenté ci-dessous.
import asyncio
from pyflink.table import DataTypes
from pyflink.table.udf import AsyncScalarFunction, udf
# Method 1: Use a Python async function
@udf(result_type=DataTypes.STRING())
async def async_api_call(product_id: str) -> str:
await asyncio.sleep(0.05)
return f"product_{product_id}"
# Method 2: Subclass the asynchronous function class
class AsyncUserLookup(AsyncScalarFunction):
def open(self, function_context):
self.cache = {}
async def eval(self, user_id: str) -> str:
if user_id in self.cache:
return self.cache[user_id]
await asyncio.sleep(0.05)
result = f"user_{user_id}"
self.cache[user_id] = result
return result
def close(self):
self.cache.clear()
async_user_lookup = udf(
AsyncUserLookup(),
input_types=[DataTypes.STRING()],
result_type=DataTypes.STRING()
)
Les fonctions définies par l'utilisateur asynchrones s'enregistrent et s'utilisent de la même manière que les fonctions synchrones.
Paramètres de configuration
Les paramètres suivants contrôlent le comportement d'exécution des fonctions définies par l'utilisateur asynchrones.
|
Paramètre |
Valeur par défaut |
Description |
|
table.exec.async-scalar.max-concurrent-operations |
10 |
Nombre maximal d'appels asynchrones simultanés par instance d'opérateur. Valeur par défaut : 10. |
|
table.exec.async-scalar.timeout |
3 min |
Délai d'expiration pour un appel asynchrone unique. |
|
table.exec.async-scalar.retry-strategy |
FIXED_DELAY |
Stratégie de nouvelle tentative après l'échec d'un appel asynchrone. Valeurs prises en charge :
|
|
table.exec.async-scalar.retry-delay |
100 ms |
Temps d'attente pour les nouvelles tentatives à délai fixe. Remarque
Effectif uniquement lorsque table.exec.async-scalar.retry-strategy est défini sur FIXED_DELAY. |
|
table.exec.async-scalar.max-attempts |
3 |
Nombre maximal de tentatives avant qu'un appel asynchrone ne soit considéré comme échoué. Remarque
Effectif uniquement lorsque table.exec.async-scalar.retry-strategy est défini sur FIXED_DELAY. |