Polars et Delta depuis une Durable Function
Un format de table transactionnel et une fonction qu’on ne contrôle pas s’accordent mal : l’une suppose un écrivain qui dure, l’autre peut être rejouée à tout moment. Ce que cela impose, et où passe la frontière.
Le besoin ne demandait pas de moteur distribué : quelques centaines de milliers de lignes par exécution, une transformation simple, une écriture dans une table du lac. Monter un cluster pour cela est un réflexe coûteux : il démarre en trois minutes, coûte pendant qu’il chauffe, et se justifie par une volumétrie qu’on n’a pas. Une bibliothèque qui travaille en mémoire dans un seul processus suffit, et la fonction sans serveur qui l’héberge coûte ce qu’elle consomme.
Une Durable Function se compose de deux natures qu’on ne peut pas mélanger. L’orchestrateur décrit l’enchaînement et il est **rejoué depuis le début** à chaque reprise : tout ce qu’il contient doit être déterministe, donc pas d’horloge, pas d’aléatoire, et surtout aucune entrée-sortie. Les activités font le travail et ne sont exécutées qu’une fois par appel réussi. Écrire dans le lac depuis l’orchestrateur produit autant d’écritures qu’il y a de rejeux, et le rejeu est le fonctionnement normal, pas l’exception.
import azure.durable_functions as df
def orchestrator(context: df.DurableOrchestrationContext):
# Sequencing and nothing else: this block is replayed on every resumption.
days = yield context.call_activity("list_days", context.get_input())
# The activities fan out; the wait itself stays deterministic.
tasks = [context.call_activity("write_day", d) for d in days]
written = yield context.task_all(tasks)
return {"partitions": len(written)}
main = df.Orchestrator.create(orchestrator)L’activité lit, transforme et écrit. Le mode d’écriture n’est pas un détail de configuration : c’est là que se décide ce qui se passe quand la même activité est appelée deux fois. En ajout, un rejeu double les lignes. En remplacement d’une partition, il les réécrit, et c’est la seule des deux qui supporte d’être rejouée.
import os
import polars as pl
STORAGE = {
"account_name": os.environ["ADLS_ACCOUNT"],
# Managed identity rather than a key: nothing to rotate, nothing to leak.
"use_azure_cli": "false",
"azure_storage_use_emulator": "false",
}
def write_day(day: str) -> str:
frame = (
pl.scan_parquet(f"abfss://raw@{os.environ['ADLS_ACCOUNT']}.dfs.core.windows.net/{day}/*.parquet")
.filter(pl.col("amount") > 0)
.group_by("store", "item")
.agg(pl.col("amount").sum().alias("revenue"))
.with_columns(pl.lit(day).alias("day"))
.collect()
)
frame.write_delta(
"abfss://refined@account.dfs.core.windows.net/sales",
mode="overwrite",
storage_options=STORAGE,
# Replace only the day's partition, not the table: that is what makes
# the activity replayable without doubling the rows.
delta_write_options={"partition_by": ["day"], "predicate": f"day = '{day}'"},
)
return dayReste la question que ce montage pose et ne résout pas : deux activités qui écrivent la même table au même instant. Le format garantit qu’une seule des deux transactions passe, l’autre échouant sur un conflit : ce qui est le bon comportement à condition que l’appelant retente. C’est pour cela que les activités sont découpées par partition plutôt que par lot : deux partitions distinctes ne se disputent rien, et le conflit devient l’exception au lieu d’être le régime normal.
Cette approche a une limite franche et il vaut mieux la nommer que la découvrir : tout doit tenir dans la mémoire d’une seule fonction. Au-delà, il faut soit découper plus finement, ce qui a une fin, soit revenir à un moteur distribué, et l’économie s’inverse. Je ne dis rien non plus de l’entretien de la table : le compactage des petits fichiers et le nettoyage des versions anciennes sont indispensables et ne se font pas depuis la fonction qui écrit. Ils demandent leur propre travail périodique, que ce montage ne fournit pas.