Écriture et soumission de jobs
Écrire un job PySpark
Un job PySpark minimal lit depuis PostgreSQL, applique une transformation et écrit le résultat en retour.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("nhic-example-job")
.getOrCreate()
)
df = (
spark.read.format("jdbc")
.option("url", "jdbc:postgresql://db:5432/warehouse")
.option("dbtable", "raw.hmis_malaria")
.option("user", "spark")
.option("password", "secret")
.load()
)
result = df.filter(df["confirmed_cases"] > 0).groupBy("facility_id").sum("confirmed_cases")
(
result.write.format("jdbc")
.option("url", "jdbc:postgresql://db:5432/warehouse")
.option("dbtable", "analytics.malaria_by_facility")
.option("user", "spark")
.option("password", "secret")
.mode("overwrite")
.save()
)
spark.stop()Soumission via l’API REST Livy
Livy accepte les soumissions de jobs sous forme de requêtes HTTP POST vers /batches.
{
"file": "s3://bucket/jobs/my_job.py",
"className": "main",
"args": ["--date", "2025-01"],
"conf": {
"spark.executor.memory": "4g"
}
}Après la soumission, interrogez GET /batches/{id}/state jusqu’à ce que l’état soit success ou dead.
Déclenchement depuis Prefect
Les workflows Prefect utilisent spark.utils.py (situé dans apps/analytics/prefect/.prefect/workflows/) pour soumettre des jobs Spark dans le cadre de pipelines plus larges. Les jobs Spark restent ainsi visibles dans l’interface Prefect — ils héritent des mêmes sémantiques de planification, de supervision et de reprise que tous les autres workflows du HIC.
Dernière mise à jour le