Nous avons préparé des guides détaillés et des exemples de scripts pour les cas d'usage courants afin de vous aider à démarrer rapidement et efficacement
Objectif
Nous avons préparé des guides détaillés et des exemples de scripts pour les cas d'usage courants afin de vous aider à démarrer rapidement et efficacement. Ces exemples montrent comment exploiter les Custom Actions dans différents scénarios, vous permettant d'extraire, transformer et charger (ETL) des données à travers les différents composants de votre data platform.
1. Custom Action avec une table du Lakehouse Manager
Aperçu :
Ce court exemple montre comment extraire des données d'une table du Lakehouse Manager, les transformer puis les insérer ou mettre à jour dans le Lakehouse Manager.
Le code est écrit dans le contexte Custom Action et utilise la table stations_rides du tutoriel Démarrage. Si vous avez suivi ce tutoriel, vous pouvez simplement copier-coller le code ci-dessous pour le tester, sinon vous devrez l'adapter à vos tables et sources de données.
Exemple d'application :
Warning
N'oubliez pas de construire la table que vous utilisez dans le Lakehouse Manager puis de la charger avec une DPE Load Action. Vous devez charger votre table avant d'utiliser le code ci-dessous, sinon cela ne fonctionnera pas.
import sys
import pandas as pd
import logging
from forepaas.dwh import connect
from forepaas.dwh import bulk_insert
logger = logging.getLogger(__name__)
def customfunc(event):
try:
logger.notice("Begin function")
# connexion au dataset par défaut
cn = connect("dwh/default_dataset/")
# option 1 : extraire les données de la table sans SQL requis
df = cn.select("stations_rides")
# option 2 : extraire les données avec une requête SQL personnalisée
df = cn.query("SELECT station_id, date, rides, station_name FROM stations_rides")
# effectuez votre transformation personnalisée dans le dataframe
df.loc[df["station_name"] == 'Harlem-Lake', "rides"] = 0
# réinsérez votre dataframe dans la table de destination
stats = bulk_insert(cn, "stations_rides", df)
# affichez les statistiques d'insertion (si compatible avec le SGBD)
logger.info(stats)
# supprimez les lignes où le nom de la station est "Davis"
cn.delete("stations_rides", {"station_name":"Davis"})
# mettez à jour les lignes en définissant rides à 0 où station_id=40040
cn.update("stations_rides", {"rides":0}, {"station_id":40040})
# une fois terminé, déconnectez cn
del cn
logger.notice("END function")
except Exception as err:
raise Exception("err:{} L:{}".format(err,sys.exc_info()[2].tb_lineno))
Aperçu :
Il arrive que vous deviez gérer un format de fichier complexe au-delà des capacités de notre Load Action.
Dans ce cas, nous vous conseillons de stocker et de manipuler les fichiers avec les Buckets de la Data Platform dans votre Projet.
Info
Dans le SDK actuel de la Data Platform, le connector Datastore est utilisé pour interagir avec les Buckets de la Data Platform. Vous pouvez simplement voir le Datastore comme un conteneur de buckets.
Exemples d'application :
import sys
import pandas as pd
from logging import getLogger
from forepaas.dwh import connect
from forepaas.dwh import bulk_insert
logger = getLogger(__name__)
def extract_func(event):
try:
# nous récupérons les données d'un bucket et nous les archiverons dans un autre bucket
bucket_source_name = "your_source_bucket_name_here"
bucket_archives_name = "your_source_bucket_name_here"
# créez un connector pour gérer le bucket
bucket_connector = connect("data_store/{}".format(bucket_source_name))
# listez les fichiers du bucket
files = bucket_connector.list()
# récupérez un fichier du bucket Data Store vers un dossier local temporaire
bucket_filepath = "stations_rides.csv"
local_filepath = "/tmp/stations_rides.csv"
bucket_connector.fget(bucket_filepath, local_filepath)
# lisez puis transformez le fichier selon vos besoins
# ici le format de la colonne date est simplement ajusté pour des raisons de compatibilité
df = pd.read_csv(local_filepath, sep=';')
df['date'] = pd.to_datetime(df['date'])
# chargez le dataframe dans une table du project nommée 'raw_file'
cn = connect("dwh/default_dataset/")
bulk_insert(cn, "stations_rides_artur", df)
del cn
# option 1 : copiez le fichier dans le bucket d'archives
bucket_archive_filepath = "archives/stations_rides.csv"
bucket_connector.fcopy_to(bucket_archives_name, bucket_archive_filepath, bucket_filepath)
# option 2 : déposez un fichier dans les archives
bucket_archives = connect("data_store/{}".format(bucket_archives_name))
bucket_archives.fput(bucket_archive_filepath, local_filepath)
del bucket_archives
# supprimez le fichier du bucket source
bucket_connector.delete(bucket_filepath)
# déconnectez du datastore
del bucket_connector
except Exception as err:
raise Exception("err:{} L:{}".format(err,sys.exc_info()[2].tb_lineno))
Voici ci-dessous un exemple de code qui upload une image vers un bucket depuis une simple URL.
from forepaas.dwh import connect
data_store = connect('data_store')
# Récupérez le bucket et uploadez l'image depuis l'URL vers le chemin uploads/test.jpg.
# Et enfin récupérez l'image depuis le bucket
bucket_test = data_store.get_bucket('test')
lists = bucket_test.list(recursive=True)
bucket_test.put_request("https://i.stack.imgur.com/r8jTK.jpg", path='uploads/test.jpg')
data = bucket.get('hello/test.jpg')
# Créez un bucket s'il n'existe pas déjà
if data_store.bucket_exists('test-exists') is False:
data_store.create_bucket('test-exists')
# Connectez-vous directement au bucket test et supprimez le fichier
bucket_test2 = connect('data_store/test')
bucket_test2.delete('hello/test.jpg')
Tip
Cela fonctionne aussi avec tout Object Store que vous pouvez définir comme source compatible S31 dans les Connectors.
3. Custom Action avec une Connectors Source
Aperçu :
Cet exemple vous montrera comment récupérer un fichier directement depuis une Connectors Source, le traiter et insérer les données traitées dans une table du Lakehouse Manager.
Le code est écrit dans le contexte Custom Action et utilise la source chicago_files du tutoriel de Démarrage. Si vous avez suivi ce tutoriel, vous pouvez simplement copier-coller le code ci-dessous pour le tester, sinon vous devrez l'adapter à vos tables et sources de données.
Tip
Cela fonctionne avec tout protocole Source tel que FTP, Dropbox, etc.
Exemples d'application :
import sys
from forepaas.dwh.connect import connect
from forepaas.dwh import bulk_insert
import logging
logger = logging.getLogger(__name__)
def customfunc(event):
try:
# ici nous nous connectons à une source nommée 'chicago_files'
source_address = "dwh/chicago_files_artur/"
# spécifiez le nom de fichier non pris en charge depuis la liste des fichiers de la source
filename_w_extension = "stations_rides.csv"
# connectez-vous au fichier pour obtenir l'adresse
source_file_connector = connect(source_address + filename_w_extension)
# connectez-vous directement au fichier
file_address = source_file_connector.get()
file_connector = connect(file_address)
# Extrayez et traitez le fichier pour qu'il soit utilisable
df = file_connector.extract(return_type='dataframe')
# Traitez les données
# - - - -
# connectez-vous au Lakehouse Manager
dm_connector = connect("dwh/default_dataset/")
# insérez dans une table de destination existante
stats = bulk_insert(dm_connector, "chicago_calendar_full", df)
logger.info(stats)
# déconnectez du datastore et de la source distante
del source_connector
del dm_connector
except Exception as err:
raise Exception(f"err:{err} L:{sys.exc_info()[2].tb_lineno}")
Aller plus loin
Si vous avez besoin d'une formation ou d'une assistance technique pour la mise en oeuvre de nos solutions, contactez votre commercial ou cliquez sur ce lien pour obtenir un devis et demander une analyse personnalisée de votre projet à nos experts de l’équipe Professional Services.
Posez vos questions, faites-nous part de vos commentaires et interagissez directement avec l’équipe qui développe la Data Platform sur le canal Discord dédié.
Si vous avez besoin d'une assistance concernant vos services OVHcloud, créez une demande depuis notre centre d'aide.
Rejoignez notre communauté d'utilisateurs.
1 : S3 est une marque déposée appartenant à Amazon Technologies, Inc. Les services de OVHcloud ne sont pas sponsorisés, approuvés, ou affiliés de quelque manière que ce soit.