MaxCompute exécute les fonctions utilisateur tabulaires (UDTF) écrites en Python à l'aide de la version 2.7. Une UDTF accepte une ligne d'entrée et renvoie zéro ou plusieurs lignes de sortie, ce qui s'avère utile pour des opérations telles que le fractionnement ou l'expansion des données.
Pour créer et utiliser une UDTF en Python 2 :
Écrivez une classe Python qui étend
BaseUDTFet implémente les méthodes requises.Enregistrez la classe en tant qu'UDTF dans MaxCompute, puis appelez-la dans MaxCompute SQL.
Structure du code UDTF
Une UDTF en Python 2 se compose de cinq composants au maximum.
| Composant | Obligatoire | Description |
|---|---|---|
| Déclaration d'encodage | Non | Déclare l'encodage du fichier. Utilisez #coding:utf-8 ou # -*- coding: utf-8 -*-. Ajoutez cette déclaration si le code contient des caractères chinois ; sans elle, MaxCompute renvoie une erreur lors de l'exécution. |
| Importations de modules | Oui | Doit inclure from odps.udf import annotate et from odps.udf import BaseUDTF. Ajoutez from odps.distcache import get_cache_file ou from odps.distcache import get_cache_table si l'UDTF fait référence à des ressources de type fichier ou table. |
| Signature de fonction | Non | Annote l'UDTF avec @annotate(<signature>) pour déclarer les types de données d'entrée et de sortie. En l'absence de signature, MaxCompute accepte tout type d'entrée mais traite toutes les valeurs de sortie comme STRING. |
| Classe dérivée | Oui | Une classe Python qui étend BaseUDTF. Cette classe contient toute la logique de l'UDTF. |
| Méthodes de classe | Oui | Implémentez au minimum process. Consultez le tableau des méthodes ci-dessous. |
Méthodes
| Méthode | Obligatoire | Description |
|---|---|---|
BaseUDTF.init() |
Non | Initialise l'état avant le traitement du premier enregistrement. Si vous redéfinissez init, appelez super(BaseUDTF, self).init() au début. Utilisez cette méthode pour configurer l'état que l'UDTF doit maintenir entre les enregistrements. |
BaseUDTF.process([args, ...]) |
Oui | Appelée une fois pour chaque ligne d'entrée. Les arguments correspondent aux paramètres d'entrée de l'UDTF tels que déclarés dans SQL. |
BaseUDTF.forward([args, ...]) |
Oui (appelée dans process) |
Émet une ligne de sortie à chaque appel. Appelez cette méthode une fois pour chaque ligne que vous souhaitez renvoyer. Si aucune signature de fonction n'est définie, convertissez tous les arguments en STRING avant d'appeler forward. |
BaseUDTF.close() |
Non | Appelée une fois avant le traitement du dernier enregistrement. Utilisez cette méthode pour libérer des ressources ou vider la sortie. |
Exemple
L'UDTF suivante fractionne une chaîne séparée par des virgules et émet chaque valeur sous forme de ligne distincte.
#coding:utf-8
from odps.udf import annotate
from odps.udf import BaseUDTF
@annotate('string -> string')
class Explode(BaseUDTF):
def process(self, arg):
props = arg.split(',')
for p in props:
self.forward(p)
Signatures de fonction et types de données
Format de signature
@annotate('arg_type_list -> type_list')
arg_type_list: liste séparée par des virgules des types de paramètres d'entrée. Utilisez*pour accepter un nombre quelconque d'arguments, ou laissez vide pour n'accepter aucun argument.type_list: liste séparée par des virgules des types de valeurs de retour. Une UDTF peut renvoyer plusieurs colonnes.
Le tableau suivant présente des exemples de signatures valides.
| Signature | Types d'entrée | Types de retour |
|---|---|---|
@annotate('bigint,boolean->string,datetime') |
BIGINT, BOOLEAN | STRING, DATETIME |
@annotate('*->string,datetime') |
Nombre quelconque d'arguments | STRING, DATETIME |
@annotate('->double,bigint,string') |
Aucun | DOUBLE, BIGINT, STRING |
@annotate("array<string>,struct<a1:bigint,b1:string>,string->map<string,bigint>,struct<b1:bigint>") |
ARRAY, STRUCT, STRING | MAP, STRUCT |
Lors de l'analyse sémantique, MaxCompute vérifie que les types de données des arguments réels correspondent à la signature. Une incompatibilité génère une erreur.
Les types de données disponibles dépendent de l'édition de types de données de votre projet MaxCompute. Pour plus d'informations, consultez Data type editions.
Mappages de types de données
Rédigez le code UDTF en utilisant les types Python correspondant aux types SQL MaxCompute.
| Type SQL MaxCompute | Type Python 2 |
|---|---|
| BIGINT | int |
| STRING | str |
| DOUBLE | float |
| BOOLEAN | bool |
| DATETIME | int (millisecondes écoulées depuis le 1er janvier 1970, 00:00:00 UTC) |
| FLOAT | float |
| CHAR | str |
| VARCHAR | str |
| BINARY | bytearray |
| DATE | int |
| DECIMAL | decimal.Decimal |
| ARRAY | list |
| MAP | dict |
| STRUCT | collections.namedtuple |
Notes complémentaires sur la gestion des types :
La valeur NULL dans MaxCompute SQL correspond à
Noneen Python.odps.udf.int(value, silent=True)renvoieNoneau lieu de lever une erreur lorsque la valeur ne peut pas être convertie en int.
Référencer des ressources de fichiers et de tables
Utilisez le module odps.distcache pour charger des ressources de type fichier ou table dans votre UDTF.
get_cache_file(resource_name): renvoie un objet semblable à un fichier pour la ressource de fichier nommée. Appelezclose()sur l'objet une fois terminé. Déclarez la ressource de fichier lors de l'enregistrement de l'UDTF ; sinon, l'appel échoue à l'exécution.get_cache_table(resource_name): renvoie un générateur pour la ressource de table nommée. Chaque itération produit un enregistrement sous forme de liste (type ARRAY).
L'exemple suivant charge un fichier JSON et une ressource de table, puis les utilise pour rechercher des ID publicitaires par ID de page.
# -*- coding: utf-8 -*-
from odps.udf import annotate
from odps.udf import BaseUDTF
from odps.distcache import get_cache_file
from odps.distcache import get_cache_table
@annotate('string -> string, bigint')
class UDTFExample(BaseUDTF):
def __init__(self):
import json
# Load the JSON file resource into a dict
cache_file = get_cache_file('test_json.txt')
self.my_dict = json.load(cache_file)
cache_file.close()
# Merge records from the table resource
records = list(get_cache_table('table_resource1'))
for record in records:
self.my_dict[record[0]] = [record[1]]
def process(self, pageid):
# Emit one row per ad ID associated with the page
for adid in self.my_dict[pageid]:
self.forward(pageid, adid)
Appeler l'UDTF dans MaxCompute SQL
Après avoir terminé le processus de développement, appelez l'UDTF depuis MaxCompute SQL :
Au sein d'un projet : appelez l'UDTF de la même manière que vous appelez les fonctions intégrées.
-
Entre projets : pour utiliser une UDTF du projet B dans le projet A, faites précéder le nom de la fonction par le nom du projet :
SELECT B:udf_in_other_project(arg0, arg1) AS res FROM table_t;Pour plus d'informations, consultez Cross-project resource access based on packages.
Limitations
MaxCompute exécute le code UDTF Python 2 dans un environnement sandbox. Les opérations suivantes ne sont pas autorisées :
Lire ou écrire dans des fichiers locaux
Démarrer des sous-processus
Démarrer des threads
Ouvrir des connexions socket
Appeler des UDF Python 2 depuis d'autres systèmes
Téléchargez uniquement du code qui utilise les bibliothèques standard Python. Les modules ou modules d'extension C qui dépendent des opérations restreintes mentionnées ci-dessus ne sont pas disponibles.
Modules d'extension C disponibles
Les modules d'extension C suivants sont disponibles dans le sandbox :
array, audioop, binascii, bisect, cmath, _codecs_cn, _codecs_hk, _codecs_iso2022, _codecs_jp, _codecs_kr, _codecs_tw, _collections, cStringIO, datetime, _functools, future_builtins, _heapq, _hashlib, itertools, _json, _locale, _lsprof, math, _md5, _multibytecodec, operator, _random, _sha256, _sha512, _sha, _struct, strop, time, unicodedata, _weakref, cPickle
Tous les modules implémentés purement en Python et ne dépendant pas de modules d'extension sont également disponibles.
Limite de taille de sortie
L'écriture dans sys.stdout ou sys.stderr est limitée à 20 Ko. Les caractères dépassant cette limite sont supprimés silencieusement.
Bibliothèques tierces
Les bibliothèques tierces, telles que NumPy, sont préinstallées dans l'environnement Python 2 de MaxCompute. L'accès aux données locales et la plupart des API d'E/S réseau sont désactivés pour les bibliothèques tierces ; seul un accès limité aux E/S réseau est disponible.