Python Komponenten-API

Eine Python-Pipelogic-Komponente ist im Kern eine process()-Funktion in src/main.py plus eine component.yml, die ihr typisiertes I/O und ihre Konfiguration beschreibt. Alles andere — Klassenform, Lifecycle-Hooks, zustandsbehaftete Muster, virtual outputs, secret-Handling — baut auf diesem einen Einstiegspunkt auf.

Diese Seite ist die vollständige Referenz der Python Komponenten-API: was process() akzeptieren und zurückgeben darf, wie zustandsbehaftete Komponenten Per-Stream-Zustand halten, wie Konfigurationswerte deinen Code erreichen und was die runtime über Imports und Threading garantiert. Für die Typsystem-Seite lies Pipelang-Typsyntax und den Typkatalog.

Für Python-Komponenten (language: py) ist src/main.py der Einstiegspunkt. Der von der runtime-Basis bereitgestellte component shim importiert dein Modul und ruft pipelogic.worker.run(your_function) auf.

Minimale Komponente

from pipelogic.worker import run, configTHRESHOLD = float(config.threshold)        # configs are read at module loaddef process(message):    return {"echoed": message}run(process)

In-Tree-Komponenten rufen run(process) auf Modul-Top-Level auf — es gibt keinen if __name__ == "__main__":-Guard. Der component shim importiert dein Modul und run(...) übernimmt die Nachrichtenschleife.

Die an run übergebene Funktion erhält ein Positionsargument pro input stream und gibt einen Wert pro output stream zurück. Die Anzahl der Argumente und die Typen werden durch worker.input_type(s) und worker.output_type(s) in der component.yml bestimmt.

config — runtime-Parameter

from pipelogic.worker import configthreshold = float(config.threshold)model_cfg = str(config.model_cfg)hidden_dims = list(config.hidden_dims)     # for [Int64] etc.

config.<name> exponiert jeden in config_schema: deklarierten Parameter. Lies ihn bei Modul-Load (auf Top-Level zur __init__-Zeit) oder innerhalb der Komponentenfunktion — beides funktioniert, aber das Lesen zur Load-Zeit schlägt schnell fehl, wenn der Parameter fehlt.

config-Werte sind typisiert basierend auf dem type: in config_schema. Caste explizit, wenn du einen eingebauten Python-Typ willst:

type: in YAMLEmpfohlener Python-Cast
Stringstr(config.x)
Boolbool(config.x)
UInt64 / Int64int(config.x)
Doublefloat(config.x)
(Double, Double, Double)tuple(float(v) for v in config.x)
[Double][float(v) for v in config.x]

Mutable Parameter

Für Parameter, die in config_schema als mutable: true deklariert sind, MUSS die Komponente config.sync() innerhalb von process() aufrufen, um Updates zu erhalten. config.sync() gibt die Menge der Parameternamen zurück, die sich geändert haben seit dem vorherigen Sync (eine leere Menge, wenn sich nichts geändert hat); nach dem Aufruf geben nachfolgende config.<name>-Lesungen die neuen Werte zurück.

from pipelogic.worker import configdef process(message):    changed = config.sync()                    # set[str] — names changed since last sync    if "threshold" in changed:        ...                                    # react to a fresh threshold value    thr = float(config.threshold)              # always the latest synced value    return {"out": ...}

Ohne config.sync() scheinen mutable Parameter auf dem beim vorherigen Sync (oder bei Modul-Load) erfassten Wert festzuhängen — selbst nachdem ppl backend change-parameter auf dem backend erfolgreich war.

Funktionssignatur — einzelner Input / einzelner Output

def process(in_msg):    out_msg = ...    return out_msg

in_msg ist ein Python-Objekt, dessen Form mit dem Input-Typ übereinstimmt:

Input-Typ-AusdruckTyp von in_msg
Stringstr
Boolbool
UInt64 / Int64int
Doublefloat
Imagepipelogic.cv.Image (nach dem Wrappen; siehe unten)
DepthImagepipelogic.cv.DepthImage (unverified-by-components — see component-api/python/image)
AudioFramepipelogic.audio.AudioFrame
Tensorpipelogic.tensor.Tensor (also pipelogic.cv.Tensor)
[T]list[T]
(T1, T2, …)tuple von T1, T2, …
{f1: T1, f2: T2}dict mit String-Keys

Der Rückgabewert folgt dem Output-Typ-Ausdruck analog.

Funktionssignatur — mehrere Inputs / Outputs

Für worker.input_types: [T1, T2] nimmt die Funktion zwei Positionsargumente:

def process(image_in, mask_in):    ...    return result

Für worker.output_types: [T1, T2] gib ein tuple zurück:

def process(image_in):    return aligned_image, metadata_string

Position zählt — Output 0 ist das erste Tuple-Element, Output 1 das zweite.

Wrappen von Plattform-Record-Typen

Input-Objekte strukturierter Typen (Image, AudioFrame, Tensor) treffen als rohe Record-Objekte ein. Um sie ergonomisch zu nutzen, wrappe sie durch die passende Helper-Klasse:

from pipelogic.cv import Image, ColorSpacedef process(in_image):    img = Image(in_image)             # accepts the raw pipelang record    bgr = img.to_bgr()                # (h, w, 3) uint8 BGR ndarray    ...    out_arr = ...                     # (h, w, 3) uint8 BGR    return Image(out_arr, ColorSpace.BGR)

Image(in_image) erkennt den eingehenden Farbraum automatisch aus dem Tag des pipelang-Records. Zum In-Place-Normalisieren nutze img.convert(ColorSpace.BGR); zum Hineinschauen ohne Mutation nutze img.to_bgr() / img.to_rgb() / img.to_gray() (jeweils ein numpy-View im angeforderten Raum).

Siehe component-api/python/image für die vollständige Image-API. Andere Wrapper leben in den geschwisterlichen Top-Level-Modulen:

WrapperImport
Imagefrom pipelogic.cv import Image
AudioFramefrom pipelogic.audio import AudioFrame
Tensorfrom pipelogic.tensor import Tensor (or from pipelogic.cv import Tensor)
Meshfrom pipelogic.mesh import Mesh
Geometryfrom pipelogic.geometry import Polygon, BoundingBox, Point

DepthImage wird ebenfalls aus pipelogic.cv exportiert, aber seine Python-API wird von keiner In-Tree-Komponente ausgeübt — behandle sie als unverified-by-components (siehe component-api/python/image).

Rückgabe generischer Werte

Für Records und Tuples gib ein dict oder tuple mit der richtigen Form zurück. Ein Record-dict trägt exakt die Feldnamen des benannten Typs, lies das Layout also zuerst in type-api/catalogBoundingBox verschachtelt ein Rectangle aus zwei Points und hat kein flaches x/y/w/h:

# output_type: "BoundingBox"return BoundingBox.from_xyxy(10.0, 20.0, 40.0, 60.0, class_id=0, confidence=0.9)# output_type: "BoundingBox" — derselbe Wert in dict-Formreturn {"class": {"id": 0, "confidence": 0.9},        "rectangle": {"top_left": {"x": 10.0, "y": 20.0},                      "bottom_right": {"x": 40.0, "y": 60.0}}}# output_type: "[BoundingBox]"return [BoundingBox.from_xyxy(10.0, 20.0, 40.0, 60.0, class_id=0, confidence=0.9),        BoundingBox.from_xyxy(50.0, 60.0, 70.0, 80.0, class_id=1, confidence=0.8)]# output_type: "(Image, String)"return img, "session-id"

Die runtime serialisiert den Rückgabewert in deinem Namen in das Plattform-Protokoll. Die Schlüssel prüft sie dabei nicht gegen den deklarierten Typ: ein dict, dessen Feldnamen nicht zum Record passen, erreicht den nächsten Vertex als fehlender Wert statt als Fehler.

Zustandsbehaftete Komponenten

Wenn deine Komponente Zustand über Nachrichten hinweg tragen muss (ein Tracker, ein laufender Durchschnitt, eine Session-Map), übergib initial_state= an run:

def process(message, state):    state["count"] = state.get("count", 0) + 1    out = {"running_count": state["count"]}    return {"output": out, "state": state}    # dict with output + state keysrun(process, initial_state={"count": 0})

Die Funktion erhält den persistierten Zustand über das state=-Schlüsselwortargument und MUSS ein Objekt zurückgeben, dessen output-Key/Attribut die Nachricht(en) für die output stream(s) hält und dessen state-Key/Attribut den neuen persistierten Zustand hält. Die Rückgabe eines bloßen Werts oder eines Tuples funktioniert nicht — die runtime schlägt rets["output"] / rets.output und rets["state"] / rets.state nach und wirft, wenn state fehlt.

Für zustandsbehaftete Komponenten mit mehreren Outputs ist output ein tuple, das der deklarierten output-stream-Anzahl entspricht. Für zustandsbehaftete virtual-output-Komponenten füge virtual_output hinzu; für den streaming-Modus füge pull_request hinzu.

Die Plattform persistiert Zustand über Container-Neustarts hinweg (saubere Crash-Recovery). Beim allerersten Start ist der Zustand der an initial_state übergebene Wert.

Reinheit und Scheduling

Verwende initial_state nicht, es sei denn, die Komponente ist tatsächlich verlaufsabhängig. Ohne persistierten Zustand ist jeder Aufruf unabhängig relativ zu den gewählten inputs, der aktuellen config/files und deterministischem Code. Das ist die Form, die die Plattform sicher parallelisieren, batchen oder skalieren kann, wenn Scheduling/Skalierung es zulässt, weil kein vorheriger tick konsultiert werden muss, bevor der nächste läuft.

Sobald initial_state gesetzt ist, ist der vertex explizit zustandsbehaftet. Die runtime muss den zurückgegebenen Zustand erhalten und ihn in den nächsten tick einspeisen, was das Ordering zu einem Teil des Vertrags macht. Nutze das für Tracker, Akkumulatoren, Session-Maps, Fenster und andere Logik, bei der derselbe Input je nach vorherigen Nachrichten einen anderen Output produzieren kann.

Modul-Globals sind weiterhin nützlich für Modell-Handles, Datenbank-Clients, kompilierte Regexes und regenerierbare Caches. Sie sollten keinen semantischen Zustand enthalten, der die Output-Historie ändert. Wenn der Zustand ändert, was die Komponente bedeutet, leg ihn in state; wenn der Code mit der Außenwelt spricht, leg diese Interaktion hinter virtual input/output.

Streaming (fortgeschritten) — StreamPullShift

Für Komponenten, die Nicht-1:1-Input/Output-Raten verarbeiten müssen, übergib initial_pull_request= an run(...). Die Funktion erhält dann Sequenzen (eine pro input stream) und gibt Sequenzen plus einen pull_request (Singular) zurück.

Ein pull request kann eine Option oder eine geordnete Liste von Optionen sein. Jede Option enthält einen StreamPullShift pro input stream. Die runtime prüft Optionen der Reihe nach und führt die Komponente mit der ersten Option aus, die aktuell erfüllbar ist.

from pipelogic.worker import run, StreamPullShiftN_INPUTS = 2Z = StreamPullShift.from_fixed(0, 0, 0)ONE = StreamPullShift.from_fixed(0, 1, 1)TWO_PLUS = StreamPullShift.from_least(0, 2)def process(*input_seqs):    # input_seqs is N_INPUTS sequences (one per input stream).    # ... compute outputs ...    return {        "output": (out_seq_0, out_seq_1),        "pull_request": [            [ONE, TWO_PLUS],            [Z, ONE],        ],    }run(    process,    initial_pull_request=[        [ONE, TWO_PLUS],        [Z, ONE],    ],)

StreamPullShift.from_fixed(preshift, size, postshift) überspringt preshift, exponiert exakt size Nachrichten an process und rückt dann um postshift vor. StreamPullShift.from_least(preshift, size) überspringt preshift und feuert dann, wenn mindestens size Nachrichten verfügbar sind. initial_pull_request= setzt den ersten Batch; danach muss jeder streaming-Rückgabewert pull_request enthalten.

Virtuelle Eingaben / Ausgaben

Virtual input/output ist die Nebeneffekt-Grenze. Nutze es, wenn eine Komponente einen Side-Thread, ein Callback, einen HTTP/WebSocket-Listener, einen Device-Reader, einen Modell-Stream oder einen externen Sender braucht, während die graph-zugewandte Funktion als scheduler-sichtbare Stream-Funktion erhalten bleibt.

Eingehende Nebeneffekte gehen durch virtual_in.push(...). Die process-Funktion empfängt sie über einen virtual_input-Schlüsselwortparameter, wenn der pull request den virtual input auswählt. Diese Payloads sind intern für die Komponente und die Bridge, keine strikten component.yml-Stream-Typen; ein virtueller Kanal kann mehrere Payload-Formen tragen, daher sollte der Komponenten-Code Payloads explizit taggen, validieren, dispatchen und fehlerhafte zurückweisen. Ausgehende Nebeneffekte gehen durch virtual_output: process gibt Work-Items zurück, und ein mit virtual_out.set_on_push(...) registriertes Bridge-Callback führt das eigentliche IO aus.

from threading import Threadfrom pipelogic.worker import run, StreamPullShift, virtual_in, virtual_outZ = StreamPullShift.from_fixed(0, 0, 0)ONE = StreamPullShift.from_fixed(0, 1, 1)def side_reader():    for event in read_external_http_or_device_events():        virtual_in.push(event, blocking=False)    virtual_in.close()Thread(target=side_reader, daemon=True).start()def process(graph_msgs, *, virtual_input):    return {        "output": transform(graph_msgs, virtual_input),        "virtual_output": build_side_effect_work_items(graph_msgs),        "pull_request": [            [ONE, Z],  # virtual input first, then declared graph input            [Z, ONE],        ],    }def send_side_effect(work_item):    deliver_to_external_system(work_item)virtual_out.set_on_push(send_side_effect)run(    process,    use_virtual_output=True,    initial_pull_request=[        [ONE, Z],        [Z, ONE],    ],)

Zwei Ordering-Regeln zählen:

  • Wenn process virtual_input deklariert, hat jede pull-Option den virtual-input-Shift zuerst, dann einen Shift pro deklariertem graph input.
  • Wenn use_virtual_output=True, muss process eine virtual_output-Liste zurückgeben. Zustandsbehaftete Komponenten geben auch state zurück; streaming-Komponenten geben auch pull_request zurück.

Was der component shim für dich erledigt

Du musst nicht:

  • Streams öffnen oder schließen.
  • (De)serialisierung typisierter Nachrichten auf der Leitung übernehmen.
  • Shared Memory zwischen vertices auf demselben Node verwalten.
  • Health-Checks oder sauberes Herunterfahren implementieren.

Du musst:

  • Die Funktion deterministisch genug machen, dass die Plattform sie beim Retry wiederholen kann (kein globaler veränderbarer Zustand außerhalb des state-Arguments).
  • print() (oder Pythons logging-Modul) für jede Komponenten-Log-Ausgabe nutzen. print() ist der kanonische Komponenten-Logging-Pfad — die runtime erfasst stdout / stderr und zeigt es über ppl deployment logs an.
  • config.sync() innerhalb von process() aufrufen, wann immer du dich auf mutable: true Parameter verlässt; andernfalls bleiben die Werte beim vorherigen Sync eingefroren.

Imports-Spickzettel

from pipelogic.worker import run, config, StreamPullShift, virtual_in, virtual_outfrom pipelogic.cv import Image, ColorSpace, Tensor, IMAGE_TYPEfrom pipelogic.audio import AudioFramefrom pipelogic.tensor import Tensor, TENSOR_TYPEfrom pipelogic.geometry import Polygon, BoundingBox, Pointfrom pipelogic.mesh import Meshfrom pipelogic.chat import ChatBackend, Message, run_chat, VisionBackend, run_visionfrom pipelogic.infer import HotSwapModel, hf_login, hf_snapshot_download, resolve_devicefrom pipelogic.cloud import PointCloudfrom pipelogic.types import (    Named, List, Union, Tuple,    register_named, get_type, create_type,)

pipelogic.types exponiert die Typsystem-Konstruktoren (Named, List, Union, Tuple) und die Registry-Helfer (register_named, get_type, create_type). Wenn ein Name nicht hier ist, prüfe type-api/catalog für den Katalog.

Häufige Fehler

  • numpy direkt zu importieren ist in Ordnung; lege es niemals in requirements.txt — die runtime-Basis pinnt numpy bereits. Es erneut zu pinnen kann ABI-Konflikte erzeugen.
  • Konstruiere Plattform-Records nicht als rohe Dicts — für Image gehe immer über pipelogic.cv.Image(...). Dasselbe für AudioFrame, Tensor usw. Handgebaute Dicts überspringen die Validierung, die buffer-vs-shape-Mismatches an der Grenze abfängt.
  • sys.exit() / os._exit() aus der Komponente vermeiden — lass Exceptions propagieren. Der shim fängt sie ab und meldet sie; explizite Exits sehen wie Crashes aus und verwirren die Retry-Semantik. (os.exit() ist keine echte Funktion — sys.exit() wirft SystemExit und os._exit() beendet den Prozess; beide sind hier falsch.)
  • config.sync() innerhalb von process() vergessen beim Lesen von mutable: true Parametern. Ohne den Aufruf bleibt der Wert beim vorherigen Sync (oder bei Modul-Load) eingefroren, selbst nachdem ppl backend change-parameter auf dem backend erfolgreich war.

Verwandt

War diese Seite hilfreich?