Batch-Import

Dieser Leitfaden baut das Gegenstück zur FastAPI-Integration: keinen Service, der auf Anfragen wartet, sondern ein Skript, das einen Datenbestand einmal abarbeitet und sich dann beendet. Es importiert Rechnungen nach enaio:

  • Ein Ordner pro Rechnungsjahr, angelegt bei der ersten Verwendung.

  • Ein Register pro Monat in diesem Ordner.

  • Ein Dokument pro Rechnung, mit der Rechnungsdatei als Anhang.

  • Jeder Schritt über einen fachlichen Schlüssel identifiziert, ein zweiter Lauf aktualisiert also statt zu duplizieren.

  • Ein fehlerhafter Datensatz wird protokolliert und gezählt, der Lauf geht weiter.

Ein Importer ist der klassische Fall für den SyncPoolClient. Die Arbeit läuft sequenziell, eine Rechnung nach der anderen, es gibt für einen Event Loop also nichts zu verschränken, und async/await würde nur Rauschen hinzufügen. Der async-Client lohnt sich dort, wo viele Operationen gleichzeitig warten, etwa in einem Web-Service.

1. Installation

  • uv

  • pip

uv add ecmind-blue-client
pip install ecmind-blue-client

Alles Weitere steckt in der Standardbibliothek: csv für die Eingabe, pathlib für die Dateien, collections.Counter für die Statistik, logging für das Protokoll.

2. Eingabedaten

Der Importer liest eine CSV, deren Pfade in der Spalte file relativ zur CSV aufgelöst werden. Daten und Dateien lassen sich damit gemeinsam verschieben:

invoice_number,invoice_date,supplier,amount,file
4711,2024-03-01,Example Supplier AG,1234.50,files/invoice-4711.txt
4712,2024-03-14,Example Supplier AG,87.00,files/invoice-4712.txt
Spalte Bedeutung

invoice_number

Fachlicher Schlüssel der Rechnung. Identifiziert das Dokument auf dem Server und muss daher eindeutig sein.

invoice_date

Rechnungsdatum im ISO-Format. Bestimmt Jahresordner und Monatsregister.

supplier

Name des Lieferanten, optional.

amount

Rechnungsbetrag, optional.

file

Pfad zur Rechnungsdatei, relativ zur CSV.

3. Modellklassen

Drei Objekttypen, einer pro Ebene der Hierarchie. In einem echten Projekt werden sie mit ecm-generate-models generiert, von Hand geschrieben sehen sie so aus:

# models.py
from datetime import date

from ecmind_blue_client.ecm.model import (
    ECMDocumentModel,
    ECMField,
    ECMFolderModel,
    ECMRegisterModel,
)


class InvoiceFolder(ECMFolderModel):
    _internal_name_ = "InvoiceFolder"

    Title: ECMField[str] = ECMField(str, mandatory=True)
    Year: ECMField[int] = ECMField(int, default=None)


class InvoiceRegister(ECMRegisterModel):
    _internal_name_ = "InvoiceRegister"

    Name: ECMField[str] = ECMField(str, mandatory=True)


class InvoiceDocument(ECMDocumentModel):
    _internal_name_ = "InvoiceDocument"

    Title: ECMField[str] = ECMField(str, mandatory=True)
    InvoiceNumber: ECMField[str] = ECMField(str, default=None)
    Supplier: ECMField[str] = ECMField(str, default=None)
    InvoiceDate: ECMField[date] = ECMField(date, default=None)
    Amount: ECMField[float] = ECMField(float, default=None)

4. Konfiguration und Start

Zugangsdaten gehören in die Umgebung, nicht in das Skript. Der Pool wird einmal erzeugt und lebt für den gesamten Lauf, und die Zeitzone der Installation wird vor dem ersten Job gesetzt, denn das Rechnungsdatum ist ein Zeitstempel, den der Server in seiner eigenen Zone berechnet:

import logging
import os

from ecmind_blue_client import set_server_timezone
from ecmind_blue_client.ecm import ECM
from ecmind_blue_client.pool import SyncPoolClient

logging.basicConfig(level=logging.INFO, format="%(levelname)-7s %(name)s: %(message)s")
log = logging.getLogger("batch-import")

set_server_timezone(os.environ.get("ECMIND_BLUE_SERVER_TIMEZONE", "Europe/Zurich"))

client = SyncPoolClient(
    servers=os.environ["ECM_SERVERS"],
    username=os.environ["ECM_USERNAME"],
    password=os.environ["ECM_PASSWORD"],
)
ecm = ECM(client)

Ein sequenzieller Importer beschäftigt eine Verbindung, der Standardwert des Pools (pool_size=10) genügt also bei Weitem. Bei einem Lauf über Stunden verhindert keepalive_interval=300, dass unbenutzte Verbindungen zwischenzeitlich von einer Firewall verworfen werden, siehe KeepAlive im Pool.

client.close() gehört in einen finally-Block. Es stoppt die Hintergrund-Worker und schliesst die unbenutzten Verbindungen, was spätestens dann zählt, wenn der Importer aus einer Zeitsteuerung läuft, die ein sauberes Beenden des Prozesses erwartet.

5. Ordner und Register pro Rechnung

Beide Ebenen werden mit einem upsert() über ihren eigenen fachlichen Schlüssel angelegt, der Importer muss also nicht wissen, ob sie bereits existieren. Die ID wird pro Lauf zwischengespeichert, damit die zweite Rechnung desselben Monats den Server nicht erneut fragt:

def folder_for_year(ecm, year: int, cache: dict[int, int]) -> int:
    """Gibt die ID des Ordners für ein Rechnungsjahr zurück und legt ihn bei Bedarf an."""
    if year in cache:
        return cache[year]

    title = f"Rechnungen {year}"
    object_id, _, _, action = (
        ecm.dms.upsert(InvoiceFolder(Title=title, Year=year))
        .search(InvoiceFolder.Title == title)
        .execute()
    )
    log.info("Ordner %r: %s (id %s)", title, action, object_id)
    cache[year] = object_id
    return object_id


def register_for_month(ecm, folder_id: int, month_key: str, cache: dict[str, int]) -> int:
    """Gibt die ID des Monatsregisters in einem Jahresordner zurück und legt es bei Bedarf an."""
    if month_key in cache:
        return cache[month_key]

    object_id, _, _, action = (
        ecm.dms.upsert(InvoiceRegister(Name=month_key), folder_id=folder_id)
        .search(InvoiceRegister.Name == month_key)
        .execute()
    )
    log.info("Register %r: %s (id %s)", month_key, action, object_id)
    cache[month_key] = object_id
    return object_id
Der Suchschlüssel muss für sich eindeutig sein

Der <Search>-Abschnitt eines Upserts wird allein aus den .search()-Bedingungen gebaut. folder_id und register_id platzieren das neue Objekt, sie schränken die Suche nicht ein.

Ein Register mit dem Namen März würde daher auch das März-Register jedes anderen Jahres treffen, und das voreingestellte .action_multiple("ERROR") würde den Datensatz abbrechen, statt ein Register im falschen Ordner zu aktualisieren. Deshalb trägt month_key das Jahr mit: 2024-03.

upsert() nimmt für folder_id und register_id die numerischen IDs. insert() und insert_and_get() akzeptieren zusätzlich eine Modellinstanz und lesen die ID daraus.

6. Das Dokument samt Datei

Das Dokument liegt im Register seines Monats und bringt die Rechnungsdatei mit. .files(…​, replace=True) bedeutet, dass ein zweiter Lauf die Datei ersetzt und keine zweite Kopie anhängt:

from datetime import date
from pathlib import Path

from ecmind_blue_client.rpc import JobRequestFileFromPath


def import_document(ecm, row: dict[str, str], base_dir: Path, folder_id: int, register_id: int) -> str:
    """Upsert eines Rechnungsdokuments samt Datei, gibt die ausgeführte Aktion zurück."""
    file_path = (base_dir / row["file"]).resolve()
    if not file_path.is_file():
        raise OSError(f"Datei nicht gefunden: {file_path}")

    document = InvoiceDocument(
        Title=f"Rechnung {row['invoice_number']}",
        InvoiceNumber=row["invoice_number"],
        Supplier=row["supplier"] or None,
        InvoiceDate=date.fromisoformat(row["invoice_date"]),
        Amount=float(row["amount"]) if row["amount"] else None,
    )

    _, _, _, action = (
        ecm.dms.upsert(document, folder_id=folder_id, register_id=register_id)
        .search(InvoiceDocument.InvoiceNumber == row["invoice_number"])
        .files([JobRequestFileFromPath(file_path)], replace=True)
        .execute()
    )
    return action

JobRequestFileFromPath liest die Datei erst beim Senden des Jobs, es wird also vorher nichts im Speicher gepuffert. Inhalt, der bereits im Speicher liegt, geht über JobRequestFileFromBytes(data, extension).

Das zurückgegebene action ist, was der Server tatsächlich getan hat, "INSERT" oder "UPDATE", und damit die ehrliche Grundlage der Laufstatistik. Wer stattdessen strikt nur anlegen will, ergänzt .action1("NONE") und überspringt bereits vorhandene Datensätze, siehe upsert().

7. Ein fehlerhafter Datensatz darf den Lauf nicht beenden

Ein Importer, der bei Datensatz 700 von 1000 stirbt, ist schlechter als einer, der 999 Erfolge und einen Fehler meldet. Jeder Datensatz läuft daher in seinem eigenen try-Block, und gefangen wird nur, was einen einzelnen Datensatz betrifft:

from ecmind_blue_client.ecm import ECMException

# ECMException deckt alles ab, was der Server ablehnt, ValueError die clientseitigen Prüfungen
# (Pflicht- und schreibgeschützte Felder) sowie ein defektes Datum oder einen defekten Betrag,
# KeyError eine fehlende CSV-Spalte und OSError eine nicht lesbare Datei.
ROW_ERRORS = (ECMException, ValueError, KeyError, OSError)
counts: Counter[str] = Counter()

for line_no, row in enumerate(csv.DictReader(handle), start=2):
    label = row.get("invoice_number") or f"Zeile {line_no}"
    try:
        invoice_date = date.fromisoformat(row["invoice_date"])
        folder_id = folder_for_year(ecm, invoice_date.year, folders)
        month_key = f"{invoice_date.year}-{invoice_date.month:02d}"
        register_id = register_for_month(ecm, folder_id, month_key, registers)

        action = import_document(ecm, row, base_dir, folder_id, register_id)
        counts[action.lower()] += 1
        log.info("Rechnung %s: %s", label, action)
    except ROW_ERRORS as error:
        counts["failed"] += 1
        log.warning("Rechnung %s fehlgeschlagen: %s", label, error)

Bewusst ungefangen bleibt alles, was den Rest des Laufs sinnlos macht: eine fehlende Umgebungsvariable, ein nicht erreichbarer Server, eine nicht vorhandene CSV. Das scheitert sofort und mit Traceback, und genau das soll eine Zeitsteuerung sehen.

Der Exit-Code macht aus der Statistik etwas, worauf ein Cron-Job oder eine Pipeline reagieren kann:

log.info(
    "fertig: %s angelegt, %s aktualisiert, %s fehlgeschlagen",
    counts["insert"],
    counts["update"],
    counts["failed"],
)
return 1 if counts["failed"] else 0

8. Ausführen

export ECM_SERVERS="enaio.example.com:4000:1"
export ECM_USERNAME="<user>"
export ECM_PASSWORD="<password>"
export ECMIND_BLUE_SERVER_TIMEZONE="Europe/Zurich"

uv run batch_import.py
INFO    batch-import: Ordner 'Rechnungen 2024': INSERT (id 4711)
INFO    batch-import: Register '2024-03': INSERT (id 4712)
INFO    batch-import: Rechnung 4711: INSERT
INFO    batch-import: Rechnung 4712: INSERT
INFO    batch-import: fertig: 2 angelegt, 0 aktualisiert, 0 fehlgeschlagen

Beim zweiten Lauf meldet jede Zeile UPDATE, und es entsteht nichts doppelt. Das ist der Sinn des fachlichen Schlüssels auf jeder Ebene, und es macht den Importer nach einem Teilabbruch gefahrlos wiederholbar.

9. Vollständiges Beispiel

Das folgende Skript liegt lauffähig im Repository als examples/batch_import.py, zusammen mit den Modellklassen in examples/models.py und den Beispieldaten in examples/invoices.csv:

"""Batch importer: file invoice files into enaio, driven by a CSV."""

import csv
import logging
import os
import sys
from collections import Counter
from datetime import date
from pathlib import Path

from ecmind_blue_client import set_server_timezone
from ecmind_blue_client.ecm import ECM, ECMException
from ecmind_blue_client.pool import SyncPoolClient
from ecmind_blue_client.rpc import JobRequestFileFromPath

from models import InvoiceDocument, InvoiceFolder, InvoiceRegister

log = logging.getLogger("batch-import")

ROW_ERRORS = (ECMException, ValueError, KeyError, OSError)


def folder_for_year(ecm, year: int, cache: dict[int, int]) -> int:
    """Return the ID of the folder for one invoice year, creating it if needed."""
    if year in cache:
        return cache[year]

    title = f"Invoices {year}"
    object_id, _, _, action = (
        ecm.dms.upsert(InvoiceFolder(Title=title, Year=year))
        .search(InvoiceFolder.Title == title)
        .execute()
    )
    log.info("folder %r: %s (id %s)", title, action, object_id)
    cache[year] = object_id
    return object_id


def register_for_month(ecm, folder_id: int, month_key: str, cache: dict[str, int]) -> int:
    """Return the ID of the month register inside a year folder, creating it if needed.

    The search section of an upsert is built from the .search() conditions alone, so month_key
    carries the year as well ("2024-03"). A bare month name would match the register of every
    other year, and the default .action_multiple("ERROR") would then abort the row.
    """
    if month_key in cache:
        return cache[month_key]

    object_id, _, _, action = (
        ecm.dms.upsert(InvoiceRegister(Name=month_key), folder_id=folder_id)
        .search(InvoiceRegister.Name == month_key)
        .execute()
    )
    log.info("register %r: %s (id %s)", month_key, action, object_id)
    cache[month_key] = object_id
    return object_id


def import_document(ecm, row: dict[str, str], base_dir: Path, folder_id: int, register_id: int) -> str:
    """Upsert one invoice document with its file and return the action the server performed."""
    file_path = (base_dir / row["file"]).resolve()
    if not file_path.is_file():
        raise OSError(f"file not found: {file_path}")

    document = InvoiceDocument(
        Title=f"Invoice {row['invoice_number']}",
        InvoiceNumber=row["invoice_number"],
        Supplier=row["supplier"] or None,
        InvoiceDate=date.fromisoformat(row["invoice_date"]),
        Amount=float(row["amount"]) if row["amount"] else None,
    )

    _, _, _, action = (
        ecm.dms.upsert(document, folder_id=folder_id, register_id=register_id)
        .search(InvoiceDocument.InvoiceNumber == row["invoice_number"])
        .files([JobRequestFileFromPath(file_path)], replace=True)
        .execute()
    )
    return action


def main() -> int:
    logging.basicConfig(level=logging.INFO, format="%(levelname)-7s %(name)s: %(message)s")

    set_server_timezone(os.environ.get("ECMIND_BLUE_SERVER_TIMEZONE", "Europe/Zurich"))

    csv_path = Path(os.environ.get("ECM_IMPORT_CSV", Path(__file__).with_name("invoices.csv")))
    base_dir = csv_path.parent

    client = SyncPoolClient(
        servers=os.environ["ECM_SERVERS"],
        username=os.environ["ECM_USERNAME"],
        password=os.environ["ECM_PASSWORD"],
    )
    ecm = ECM(client)

    counts: Counter[str] = Counter()
    folders: dict[int, int] = {}
    registers: dict[str, int] = {}

    try:
        with csv_path.open(newline="", encoding="utf-8") as handle:
            # start=2 so the number in a log line matches the line in the file, header included
            for line_no, row in enumerate(csv.DictReader(handle), start=2):
                label = row.get("invoice_number") or f"line {line_no}"
                try:
                    invoice_date = date.fromisoformat(row["invoice_date"])
                    folder_id = folder_for_year(ecm, invoice_date.year, folders)
                    month_key = f"{invoice_date.year}-{invoice_date.month:02d}"
                    register_id = register_for_month(ecm, folder_id, month_key, registers)

                    action = import_document(ecm, row, base_dir, folder_id, register_id)
                    counts[action.lower()] += 1
                    log.info("invoice %s: %s", label, action)
                except ROW_ERRORS as error:
                    # One bad row must not end the run: log it, count it, keep going.
                    counts["failed"] += 1
                    log.warning("invoice %s failed: %s", label, error)
    finally:
        client.close()

    log.info(
        "done: %s inserted, %s updated, %s failed",
        counts["insert"],
        counts["update"],
        counts["failed"],
    )
    return 1 if counts["failed"] else 0


if __name__ == "__main__":
    sys.exit(main())

10. Verwandte Themen