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 |
|---|---|
|
Fachlicher Schlüssel der Rechnung. Identifiziert das Dokument auf dem Server und muss daher eindeutig sein. |
|
Rechnungsdatum im ISO-Format. Bestimmt Jahresordner und Monatsregister. |
|
Name des Lieferanten, optional. |
|
Rechnungsbetrag, optional. |
|
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.
|
|
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 Ein Register mit dem Namen |
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
-
Schnellstart für die einzelnen Schritte, die dieser Importer kombiniert
-
upsert() für Aktionssteuerung, Dateien und Tabellenfelder
-
ecm-generate-models für das Generieren der Modellklassen
-
Zeitzone der Installation für Datums- und Zeitwerte
-
KeepAlive im Pool für lang laufende Importe
-
Servererreichbarkeit prüfen für eine Prüfung vor dem ersten Job
-
FastAPI-Integration für das asynchrone Gegenstück