Tous les produits
Search
Centre de documentation

MaxCompute:Exemple : Lire des ressources MaxCompute à l'aide d'une UDTF Python 3

Dernière mise à jour :Aug 10, 2026

Utilisez une fonction de table définie par l'utilisateur (UDTF) en Python 3 pour lire des ressources MaxCompute, notamment des fichiers et des tables de ressource, depuis le client MaxCompute.

Prérequis

Avant de commencer, vérifiez que vous disposez des éléments suivants :

Méthodes de la classe Handler

Une UDTF Python 3 s'implémente sous la forme d'une classe qui étend BaseUDTF. Voici les méthodes de cette classe :

Méthode Obligatoire Description
__init__ Non Initialise l'état. Utilisez-la pour charger les fichiers et les tables de ressource une seule fois au démarrage, avant tout traitement de ligne.
process Oui Traite chaque ligne d'entrée. Appelez forward() dans cette méthode pour émettre les lignes de sortie.

Le décorateur @annotate définit la signature de la fonction. Dans la méthode process , appelez forward(*args) pour émettre chaque ligne de sortie : un appel à forward produit une seule ligne de sortie.

Paramètres dynamiques

Pour obtenir une référence complète sur les formats de signature de fonction et les types de données, consultez la section Signatures de fonction et types de données.

Selon leur position dans les signatures @annotate, les astérisques (*) ont des significations différentes :

Dans les listes de paramètres : L'astérisque * indique que l'entrée accepte un nombre indéfini de paramètres supplémentaires de tout type. Par exemple, @annotate('double,*->string') déclare un premier paramètre de type DOUBLE suivi d'un nombre variable de paramètres. Vous devez compiler le code pour calculer le nombre et les types des paramètres d'entrée, et les gérer selon la fonction printf du langage C. La méthode process reçoit ces paramètres sous la forme *args et doit les traiter explicitement.

Dans les valeurs de retour : L'astérisque * signifie qu'un nombre quelconque de valeurs STRING est renvoyé, la quantité étant déterminée par le nombre d'alias fournis lors de l'appel. Par exemple, @annotate("bigint,string->double,*") appelé sous la forme UDTF(x, y) as (a, b, c) renvoie trois valeurs : a en tant que DOUBLE, et b et c en tant que STRING. L'appel à forward dans la méthode process doit émettre un tableau dont la longueur correspond exactement au nombre d'alias ; toute divergence entraîne une erreur d'exécution, et non une erreur de compilation.

Remarque

Cette règle concernant l'astérisque * dans les valeurs de retour s'applique uniquement aux UDTF. Les fonctions d'agrégation définies par l'utilisateur (UDAF) renvoient toujours une seule valeur.

Exemple de code

Lire des ressources depuis MaxCompute

L'UDTF suivante lit les mappages page-vers-publicité à partir d'un fichier de ressource JSON (test_json.txt) et d'une table de ressource (table_resource1), puis émet une ligne par ID de publicité pour chaque ID de page d'entrée.

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):
    """Read pageid and adid_list from the file get_cache_file and the table get_cache_table to generate dict.
    """
    def __init__(self):
        import json
        cache_file = get_cache_file('test_json.txt')
        self.my_dict = json.load(cache_file)
        cache_file.close()
        records = list(get_cache_table('table_resource1'))
        for record in records:
            self.my_dict[record[0]] = record[1]
    """Enter pageid and generate pageid and all adid values.
    """
    def process(self, pageid):
        for adid in self.my_dict[pageid]:
            self.forward(pageid, adid)

Les fonctions get_cache_file et get_cache_table (issues du module odps.distcache) chargent les ressources en utilisant le nom enregistré lors de l'ajout via les commandes add file ou add table. Les ressources sont chargées une seule fois dans __init__ et réutilisées pour tous les appels à process.

Utiliser des paramètres dynamiques

L'UDTF suivante analyse une chaîne JSON et extrait les valeurs par clé. Le nombre de valeurs de retour est égal au nombre de paramètres d'entrée.

from odps.udf import annotate
from odps.udf import BaseUDTF
import json

@annotate('string,*->string,*')
class JsonTuple(BaseUDTF):
    def process(self, *args):
        length = len(args)
        result = [None] * length
        try:
            obj = json.loads(args[0])
            for i in range(1, length):
                result[i] = str(obj.get(args[i]))
        except Exception as err:
            result[0] = str(err)
            for i in range(1, length):
                result[i] = None
        self.forward(*result)

Le premier argument (args[0]) correspond à la chaîne JSON. Les arguments supplémentaires représentent les clés à extraire. La première valeur de retour contient toute erreur d'analyse, tandis que les valeurs suivantes contiennent le contenu extrait, classé par ordre de clé.

Appelez cette UDTF avec un nombre d'alias correspondant :

-- Number of output aliases matches number of input parameters
SELECT my_json_tuple(json, 'a', 'b') as (exceptions, a, b) FROM jsons;

-- The variable-length part can have no columns
SELECT my_json_tuple(json) as exceptions FROM jsons;

-- This causes a runtime error: alias count (4) does not match input parameter count (3)
SELECT my_json_tuple(json, 'a', 'b') as (exceptions, a, b, c) FROM jsons;

Enregistrer et appeler une UDTF

La procédure suivante utilise UDTFExample pour illustrer le workflow complet. Enregistrez le code sous le nom py_udtf_example.py dans le dossier bin du client MaxCompute avant de commencer.

Étape 1 : Créer des tables de ressource et préparer les données

Connectez-vous au client MaxCompute. Pour plus d'informations, consultez la page Démarrer le client MaxCompute.

Créez et alimentez la table de ressource table_resource1 :

create table if not exists table_resource1 (pageid string, adid_list array<int>);
insert into table table_resource1 values("contact_page2",array(2,3,4)),("contact_page3",array(5,6,7));
Remarque

Le champ adid_list est de type ARRAY. Pour permettre à Python 3 de lire les données de type ARRAY, exécutez la commande set odps.sql.python.version=cp37; au niveau de la session avant d'effectuer la requête.

Créez et alimentez la table interne tmp1 :

create table if not exists tmp1 (pageid string);
insert into table tmp1 values ("front_page"),("contact_page1"),("contact_page3");

Placez le fichier test_json.txt dans le dossier bin du client MaxCompute. Le fichier contient :

{"front_page":[1, 2, 3], "contact_page1":[3, 4, 5]}

Étape 2 : Enregistrer les ressources

Ajoutez le fichier Python, le fichier JSON et la table en tant que ressources MaxCompute. Pour plus d'informations, consultez la page Ajouter des ressources.

add py py_udtf_example.py;
add file test_json.txt;
add table table_resource1 as table_resource1;

Étape 3 : Créer l'UDTF

Enregistrez l'UDTF en répertoriant toutes les ressources dont elle dépend. Pour plus d'informations, consultez la page Créer une UDF.

create function my_udtf as 'py_udtf_example.UDTFExample' using 'py_udtf_example.py, test_json.txt, table_resource1';

Étape 4 : Appeler l'UDTF

Trois modèles d'appel sont pris en charge. Tous produisent la même sortie pour cet exemple.

Appel direct :

select my_udtf(pageid) as (pageid, adid) from tmp1;

Avec LATERAL VIEW :

select pageid, adid from tmp1 lateral view my_udtf(pageid) adTable as udtf_pageid, adid;

Avec LATERAL VIEW et une fonction d'agrégation :

select adid, count(1) as cnt
    from tmp1 lateral view my_udtf(pageid) adTable as udtf_pageid, adid
group by adid;

L'appel direct et les requêtes LATERAL VIEW renvoient les mêmes 9 lignes :

+------------+------------+
| pageid     | adid       |
+------------+------------+
| front_page | 1          |
| front_page | 2          |
| front_page | 3          |
| contact_page1 | 3       |
| contact_page1 | 4       |
| contact_page1 | 5       |
| contact_page3 | 5       |
| contact_page3 | 6       |
| contact_page3 | 7       |
+------------+------------+

La requête d'agrégation renvoie :

+------------+------------+
| adid       | cnt        |
+------------+------------+
| 1          | 1          |
| 2          | 1          |
| 3          | 2          |
| 4          | 1          |
| 5          | 2          |
| 6          | 1          |
| 7          | 1          |
+------------+------------+