← points de vue
13 février 2024

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.

01L’orchestrateur ne fait rien, et c’est la règle

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.

python
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)
02L’écriture, et la contrainte qu’elle porte

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.

python
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 day

Reste 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.

03Ce que je laisse de côté

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.