PyODPS vous permet de créer, lire, écrire et supprimer des tables MaxCompute par programmation. Vous pouvez également gérer les schémas, les partitions et les transferts de données en masse via MaxCompute Tunnel.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Un projet MaxCompute
PyODPS installé et configuré
Un ID AccessKey et un secret AccessKey
Un objet d'entrée MaxCompute initialisé (
o)
Démarrage rapide
Créez une table, écrivez-y des données, puis lisez-les :
from odps.models import Schema
# Define columns
schema = Schema.from_lists(['id', 'name'], ['bigint', 'string'])
# Create the table
table = o.create_table('my_table', schema, if_not_exists=True)
# Write records
o.write_table('my_table', [[1, 'Alice'], [2, 'Bob']])
# Read records
for record in o.read_table('my_table'):
print(record)
Résumé des méthodes API
|
Opération |
Méthode |
Paramètres clés |
Description |
|
Lister les tables |
|
|
Répertorie toutes les tables d'un projet |
|
Vérifier l'existence |
|
|
Vérifie si une table existe |
|
Obtenir une table |
|
|
Obtient un objet table |
|
Créer une table |
|
|
Crée une table |
|
Écrire des données |
|
|
Ajoute des enregistrements à une table |
|
Lire des données |
|
|
Lit les enregistrements d'une table |
|
Supprimer une table |
|
|
Supprime une table |
|
Convertir en DataFrame |
|
— |
Convertit une table en DataFrame PyODPS |
|
Synchroniser les métadonnées |
|
— |
Actualise l'objet table local depuis le serveur |
Pour la référence complète des méthodes PyODPS, consultez Présentation du SDK Python.
Lister les tables
Listez toutes les tables d'un projet :
for table in o.list_tables():
print(table)
Filtrez par préfixe de nom :
for table in o.list_tables(prefix="table_prefix"):
print(table.name)
Par défaut, list_tables() renvoie uniquement les noms des tables. L'accès aux propriétés telles que table_schema ou creation_time déclenche des requêtes supplémentaires et augmente la latence. À partir de PyODPS 0.11.5, passez extended=True pour récupérer ces propriétés en un seul appel :
for table in o.list_tables(extended=True):
print(table.name, table.creation_time)
Filtrez par type de table :
# Valid types: managed_table, external_table, virtual_view, materialized_view
managed_tables = list(o.list_tables(type="managed_table"))
external_tables = list(o.list_tables(type="external_table"))
virtual_views = list(o.list_tables(type="virtual_view"))
materialized_views = list(o.list_tables(type="materialized_view"))
Vérifier si une table existe
print(o.exist_table('pyodps_iris'))
# Returns True if the table exists
Obtenir les informations d'une table
Obtenez un objet table avec get_table() :
t = o.get_table('pyodps_iris')
Affichez le schéma de la table :
print(t.schema)
Sortie :
odps.Schema {
sepallength double # Sepal length (cm)
sepalwidth double # Sepal width (cm)
petallength double # Petal length (cm)
petalwidth double # Petal width (cm)
name string # Type
}
Accédez aux détails des colonnes :
# All columns
print(t.schema.columns)
# A specific column
print(t.schema['sepallength'])
# Column comment
print(t.schema['sepallength'].comment)
Accédez aux propriétés de la table :
print(t.lifecycle) # Table lifecycle
print(t.creation_time) # Creation time
print(t.is_virtual_view) # Whether the table is a view
print(t.size) # Table size in bytes
print(t.comment) # Table comment
Accéder aux tables entre projets
Passez le paramètre project pour obtenir une table d'un autre projet :
t = o.get_table('table_name', project='other_project')
Créer un schéma de table
Deux méthodes sont disponibles pour créer des schémas.
Méthode 1 : Objets Column et Partition
Utilisez les objets Column et Partition de odps.models lorsque vous avez besoin d'un contrôle total sur les définitions de colonne, y compris les commentaires :
from odps.models import Schema, Column, Partition
columns = [
Column(name='num', type='bigint', comment='the column'),
Column(name='num2', type='double', comment='the column2'),
]
partitions = [Partition(name='pt', type='string', comment='the partition')]
schema = Schema(columns=columns, partitions=partitions)
Accédez aux propriétés du schéma :
# All columns including partition columns
print(schema.columns)
# Partition columns only
print(schema.partitions)
# Non-partition column names
print(schema.names)
# Non-partition column types
print(schema.types)
Méthode 2 : Schema.from_lists()
Schema.from_lists() est plus simple, mais ne prend pas en charge les commentaires de colonne :
from odps.models import Schema
schema = Schema.from_lists(
['num', 'num2'], # Column names
['bigint', 'double'], # Column types
['pt'], # Partition names
['string'] # Partition types
)
print(schema.columns)
Créer une table
À partir d'un objet schéma
from odps.models import Schema
schema = Schema.from_lists(['num', 'num2'], ['bigint', 'double'], ['pt'], ['string'])
# Create a table
table = o.create_table('my_new_table', schema)
# Skip creation if the table already exists
table = o.create_table('my_new_table', schema, if_not_exists=True)
# Set the lifecycle (days before auto-deletion)
table = o.create_table('my_new_table', schema, lifecycle=7)
À partir de définitions de colonnes sous forme de chaînes
# Partitioned table (columns, partition columns)
table = o.create_table('my_new_table', ('num bigint, num2 double', 'pt string'), if_not_exists=True)
# Non-partitioned table
table = o.create_table('my_new_table02', 'num bigint, num2 double', if_not_exists=True)
Activer les types de données étendus
Par défaut, seuls ces types de données sont pris en charge : BIGINT, DOUBLE, DECIMAL, STRING, DATETIME, BOOLEAN, MAP et ARRAY.
Pour utiliser des types étendus tels que TINYINT et STRUCT, activez l'extension de type de données MaxCompute V2.0 :
from odps import options
options.sql.use_odps2_extension = True
table = o.create_table('my_new_table', 'cat smallint, content struct<title:varchar(100), body:string>')
Synchroniser les mises à jour de table
Lorsqu'un autre programme modifie une table, appelez reload() pour actualiser l'objet local avec les dernières métadonnées du serveur :
from odps.models import Schema
schema = Schema.from_lists(['num', 'num2'], ['bigint', 'double'], ['pt'], ['string'])
table = o.create_table('my_new_table', schema)
# Fetch the latest table metadata from the server
table.reload()
Écrire des données dans une table
write_table()
Utilisez write_table() pour des écritures simples et ponctuelles lorsque tous les enregistrements sont prêts. Cette méthode ajoute les données à la table et gère les partitions en un seul appel.
records = [
[111, 1.0],
[222, 2.0],
[333, 3.0],
[444, 4.0]
]
# Write to a partition; create the partition if it does not exist
o.write_table('my_new_table', records, partition='pt=test', create_partition=True)
Chaque appel à write_table() crée un fichier sur le serveur. Cette opération est longue et un trop grand nombre de petits fichiers dégrade les performances des requêtes. Écrivez plusieurs enregistrements par appel ou passez un objet générateur.
write_table() ajoute toujours des données. Pour remplacer les données existantes :
Tables non partitionnées : Appelez
table.truncate().Tables partitionnées : Supprimez et recréez la partition.
open_writer()
Utilisez open_writer() pour des écritures en flux continu ou pour écrire des enregistrements de manière incrémentielle au sein d'une session gérée.
t = o.get_table('my_new_table')
with t.open_writer(partition='pt=test02', create_partition=True) as writer:
records = [
[1, 1.0],
[2, 2.0],
[3, 3.0],
[4, 4.0]
]
writer.write(records) # Accepts any iterable
Écrivez dans des partitions multiniveaux :
t = o.get_table('test_table')
with t.open_writer(partition='pt1=test1,pt2=test2') as writer:
records = [
t.new_record([111, 'aaa', True]),
t.new_record([222, 'bbb', False]),
t.new_record([333, 'ccc', True]),
t.new_record([444, 'Chinese', False])
]
writer.write(records)
Écriture parallèle multiprocessus
Plusieurs processus peuvent écrire simultanément dans la même table en partageant un ID de session et en écrivant dans des blocs distincts. Chaque bloc correspond à un fichier sur le serveur. Le processus principal valide les données une fois que tous les workers ont terminé.
import random
from multiprocessing import Pool
from odps.tunnel import TableTunnel
def write_records(tunnel, table, session_id, block_id):
# Reuse the existing session
local_session = tunnel.create_upload_session(table.name, upload_id=session_id)
# Write to this process's block
with local_session.open_record_writer(block_id) as writer:
for i in range(5):
record = table.new_record([random.randint(1, 100), random.random()])
writer.write(record)
if __name__ == '__main__':
N_WORKERS = 3
table = o.create_table('my_new_table', 'num bigint, num2 double', if_not_exists=True)
tunnel = TableTunnel(o)
upload_session = tunnel.create_upload_session(table.name)
# Share the session ID across processes
session_id = upload_session.id
pool = Pool(processes=N_WORKERS)
futures = []
block_ids = []
for i in range(N_WORKERS):
futures.append(pool.apply_async(write_records, (tunnel, table, session_id, i)))
block_ids.append(i)
[f.get() for f in futures]
# Commit all blocks
upload_session.commit(block_ids)
Lire des données depuis une table
read_table()
Utilisez read_table() pour parcourir tous les enregistrements d'une table ou d'une partition :
for record in o.read_table('my_new_table', partition='pt=test'):
print(record)
head()
Prévisualisez jusqu'à 10 000 enregistrements sans effectuer une analyse complète de la table :
t = o.get_table('my_new_table')
for record in t.head(3):
print(record)
open_reader()
Utilisez open_reader() lorsque vous avez besoin d'un accès par tranches ou d'un comptage des enregistrements avant la lecture. Cette méthode expose un attribut count et prend en charge le découpage par index :
t = o.get_table('my_new_table')
with t.open_reader(partition='pt=test') as reader:
count = reader.count
for record in reader[5:10]: # Read a slice of records
print(record)
Sans bloc with :
reader = t.open_reader(partition='pt=test')
count = reader.count
for record in reader[5:10]:
print(record)
Supprimer une table
# Delete only if the table exists
o.delete_table('my_table_name', if_exists=True)
# Or call drop() on a table object
t.drop()
Convertir une table en DataFrame
to_df() convertit une table en DataFrame PyODPS. Pour plus de détails, consultez DataFrame (non recommandé).
table = o.get_table('my_table_name')
df = table.to_df()
Gérer les partitions
Vérifier si une table est partitionnée
table = o.get_table('my_new_table')
if table.schema.partitions:
print('Table %s is partitioned.' % table.name)
Parcourir les partitions
table = o.get_table('my_new_table')
# All partitions
for partition in table.partitions:
print(partition.name)
# Sub-partitions under pt=test
for partition in table.iterate_partitions(spec='pt=test'):
print(partition.name)
# Partitions matching a condition (PyODPS 0.11.3 and later)
for partition in table.iterate_partitions(spec='dt>20230119'):
print(partition.name)
À partir de PyODPS 0.11.3, iterate_partitions() accepte des expressions logiques telles que dt>20230119.
Vérifier si une partition existe
table = o.get_table('my_new_table')
table.exist_partition('pt=test,sub=2015')
Obtenir les informations d'une partition
table = o.get_table('my_new_table')
partition = table.get_partition('pt=test')
print(partition.creation_time)
print(partition.size)
Créer une partition
t = o.get_table('my_new_table')
t.create_partition('pt=test', if_not_exists=True)
Supprimer une partition
t = o.get_table('my_new_table')
t.delete_partition('pt=test', if_exists=True)
# Or call drop() on a partition object
partition.drop()
Enregistrements et mappages de types de données
Un enregistrement représente une seule ligne dans une table MaxCompute. Les quatre méthodes d'E/S — open_reader(), open_writer(), open_record_reader() et open_record_writer() — utilisent des enregistrements.
Créez un enregistrement en appelant new_record() sur un objet table.
Compte tenu de ce schéma de table :
odps.Schema {
c_int_a bigint
c_string_a string
c_bool_a boolean
c_datetime_a datetime
c_array_a array<string>
c_map_a map<bigint,string>
c_struct_a struct<a:bigint,b:string>
}
Créez et manipulez des enregistrements :
import datetime
t = o.get_table('mytable') # o is the MaxCompute entry object
# Create a record with initial values
# The number of values must match the number of fields in the schema
r = t.new_record([1024, 'val1', False, datetime.datetime.now(), None, None])
# Create an empty record
r2 = t.new_record()
# Set values by index
r2[0] = 1024
# Set values by field name
r2['c_string_a'] = 'val1'
# Set values by attribute
r2.c_string_a = 'val1'
# Set ARRAY value
r2.c_array_a = ['val1', 'val2']
# Set MAP value
r2.c_map_a = {1: 'val1'}
# Set STRUCT value (PyODPS 0.11.5 and later)
r2.c_struct_a = (1, 'val1') # tuple
r2.c_struct_a = {"a": 1, "b": 'val1'} # dict
# Get values
print(r[0]) # By index
print(r['c_string_a']) # By field name
print(r.c_string_a) # By attribute
print(r[0: 3]) # Slice
print(r[0, 2, 3]) # Multiple indices
print(r['c_int_a', 'c_double_a']) # Multiple field names
Mappages de types de données
|
Type MaxCompute |
Type Python |
|
TINYINT, SMALLINT, INT, BIGINT |
int |
|
FLOAT, DOUBLE |
float |
|
STRING |
str |
|
BINARY |
bytes |
|
DATETIME |
datetime.datetime |
|
DATE |
datetime.date |
|
BOOLEAN |
bool |
|
DECIMAL |
decimal.Decimal |
|
MAP |
dict |
|
ARRAY |
list |
|
STRUCT |
tuple / namedtuple |
|
TIMESTAMP |
pandas.Timestamp |
|
TIMESTAMP_NTZ |
pandas.Timestamp |
|
INTERVAL_DAY_TIME |
pandas.Timedelta |
Gestion des STRING
Par défaut, STRING correspond aux chaînes Unicode (str dans Python 3, unicode dans Python 2). Pour stocker des données binaires dans une colonne STRING, définissez options.tunnel.string_as_binary = True.
Comportement du fuseau horaire
PyODPS utilise le fuseau horaire local par défaut. MaxCompute ne stocke pas les valeurs de fuseau horaire ; il convertit les valeurs datetime en horodatages UNIX pour le stockage.
Pour utiliser UTC :
options.local_timezone = False
Pour utiliser un fuseau horaire spécifique :
options.local_timezone = 'Asia/Shanghai'
DECIMAL dans Python 2
Lorsque le package cdecimal est installé, PyODPS utilise cdecimal.Decimal au lieu de decimal.Decimal dans Python 2.
Comportement du type STRUCT
Avant PyODPS 0.11.5, STRUCT correspond à dict. À partir de PyODPS 0.11.5, STRUCT correspond à namedtuple par défaut.
Pour restaurer le comportement dict précédent :
options.struct_as_dict = True
Dans les environnements DataWorks, struct_as_dict est défini sur False par défaut pour des raisons de compatibilité historique. PyODPS 0.11.5 et versions ultérieures acceptent à la fois dict et tuple pour les valeurs STRUCT. Les versions antérieures n'acceptent que dict.
MaxCompute Tunnel
MaxCompute Tunnel est le canal de données bas niveau pour les chargements et téléchargements en masse. Pour la plupart des cas d'utilisation, write_table() et read_table() sont plus simples. Utilisez Tunnel lorsque vous avez besoin d'écritures multiprocessus ou d'un contrôle précis des sessions.
PyODPS ne prend pas en charge le chargement de données via des tables externes (par exemple, des tables sauvegardées par OSS ou Tablestore).
Si CPython est installé, PyODPS compile le code C lors de l'installation pour accélérer les transferts basés sur Tunnel.
Charger des données
from odps.tunnel import TableTunnel
table = o.get_table('my_table')
tunnel = TableTunnel(o)
upload_session = tunnel.create_upload_session(table.name, partition_spec='pt=test')
with upload_session.open_record_writer(0) as writer:
record = table.new_record()
record[0] = 'test1'
record[1] = 'id1'
writer.write(record)
record = table.new_record(['test2', 'id2'])
writer.write(record)
# Commit outside the with block. Committing before all data is written causes an error.
upload_session.commit([0])
Télécharger des données
from odps.tunnel import TableTunnel
tunnel = TableTunnel(o)
download_session = tunnel.create_download_session('my_table', partition_spec='pt=test')
with download_session.open_record_reader(0, download_session.count) as reader:
for record in reader:
print(record)