Python bileşen API'si

Bir Python Pipelogic komponenti, özünde src/main.py'deki bir process() fonksiyonu artı onun tipli I/O ve yapılandırmasını tanımlayan bir component.yml'dir. Geri kalan her şey — sınıf şekli, lifecycle kancaları, durum bilgili desenler, virtual output'lar, secret işleme — o tek giriş noktasının üzerine inşa edilir.

Bu sayfa tam Python bileşen API'si referansıdır: process()'in ne kabul edip döndürebileceği, durum bilgili komponentlerin stream başına durumu nasıl tuttuğu, yapılandırma değerlerinin kodunuza nasıl ulaştığı ve runtime'ın içe aktarmalar ve threading hakkında ne garanti ettiği. Tip sistemi tarafı için Pipelang tip sözdizimi ve tip kataloğu okuyun.

Python komponentleri (language: py) için src/main.py giriş noktasıdır. Runtime tabanı tarafından sağlanan component shim, modülünüzü içe aktarır ve pipelogic.worker.run(your_function)'i çağırır.

Minimal bileşen

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

In-tree komponentler run(process)'i modül üst seviyesinde çağırır — if __name__ == "__main__": koruması yoktur. Component shim modülünüzü içe aktarır ve run(...) mesaj döngüsünü devralır.

run'a geçirilen fonksiyon, input stream başına bir konumsal argüman alır ve output stream başına bir değer döndürür. Argüman sayısı ve tipler, component.yml'deki worker.input_type(s) ve worker.output_type(s) tarafından belirlenir.

config — runtime parametreleri

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>, config_schema:'da bildirilen her parametreyi sunar. Onu modül yüklemede (üst seviye __init__ zamanında) veya komponent fonksiyonu içinde okuyun — ikisi de çalışır, ancak yükleme zamanında okumak, parametre eksikse hızlı başarısız olur.

config değerleri, config_schema'daki type:'a göre tipliddir. Yerleşik bir Python tipi istiyorsanız açıkça cast'leyin:

YAML'deki type:Önerilen Python cast'i
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 parametreler

config_schema'da mutable: true olarak bildirilen parametreler için, komponent güncellemeleri almak üzere process() içinde config.sync()'i çağırMALIdır. config.sync(), önceki sync'ten bu yana değişen parametre adları kümesini döndürür (hiçbir şey değişmediyse boş bir küme); çağrıdan sonra, sonraki config.<name> okumaları yeni değerleri döndürür.

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": ...}

config.sync() olmadan, mutable parametreler önceki sync'te (veya modül yüklemede) yakalanan değerde takılı görünür — ppl backend change-parameter backend'de başarılı olduktan sonra bile.

Fonksiyon imzası — tek girdi / tek çıktı

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

in_msg, şekli girdi tipiyle eşleşen bir Python nesnesidir:

Girdi tipi ifadesiin_msg tipi
Stringstr
Boolbool
UInt64 / Int64int
Doublefloat
Imagepipelogic.cv.Image (sarmaladıktan sonra; aşağıya bakın)
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, …)T1, T2, … tuple'ı
{f1: T1, f2: T2}string key'lere sahip dict

Dönüş değeri, çıktı tipi ifadesini benzer şekilde izler.

Fonksiyon imzası — birden çok girdi / çıktı

worker.input_types: [T1, T2] için fonksiyon iki konumsal argüman alır:

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

worker.output_types: [T1, T2] için bir tuple döndürün:

def process(image_in):    return aligned_image, metadata_string

Konum önemlidir — çıktı 0 ilk tuple öğesidir, çıktı 1 ikincisidir.

Platform record tiplerini sarmalama

Yapılandırılmış tiplerin (Image, AudioFrame, Tensor) girdi nesneleri ham record nesneleri olarak gelir. Onları ergonomik şekilde kullanmak için eşleşen yardımcı sınıf aracılığıyla sarmalayın:

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), gelen renk uzayını pipelang record'unun tag'inden otomatik algılar. Yerinde normalleştirmek için img.convert(ColorSpace.BGR) kullanın; mutasyon olmadan bakmak için img.to_bgr() / img.to_rgb() / img.to_gray() kullanın (her biri istenen uzayda bir numpy view döndürür).

Tam Image API'si için component-api/python/image bakın. Diğer wrapper'lar kardeş üst seviye modüllerde yaşar:

Wrapperİçe aktarma
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 da pipelogic.cv'den dışa aktarılır, ancak Python API'si hiçbir in-tree komponent tarafından kullanılmaz — onu unverified-by-components olarak ele alın (component-api/python/image bakın).

Generic değerler döndürme

Record'lar ve tuple'lar için doğru şekle sahip bir dict veya tuple döndürün. Bir record dict'i, adlandırılmış tipin alan adlarını birebir taşır; bu yüzden önce yerleşimi type-api/catalog içinde okuyun — BoundingBox, iki Point'ten oluşan bir Rectangle içerir ve düz bir x/y/w/h taşımaz:

# output_type: "BoundingBox"return BoundingBox.from_xyxy(10.0, 20.0, 40.0, 60.0, class_id=0, confidence=0.9)# output_type: "BoundingBox" — dict biçiminde aynı değerreturn {"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"

Runtime, dönüş değerini sizin adınıza platform protokolüne serileştirir. Anahtarları bildirilen tiple karşılaştırmaz: alan adları record ile eşleşmeyen bir dict, bir sonraki vertex'e hata olarak değil, eksik bir değer olarak ulaşır.

Durum bilgili bileşenler

Komponentinizin mesajlar arasında durum taşıması gerekiyorsa (bir tracker, çalışan bir ortalama, bir session haritası), run'a initial_state= geçin:

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})

Fonksiyon, kalıcı durumu state= anahtar kelime argümanı aracılığıyla alır ve output key'i/attr'ı output stream(ler) için mesaj(lar)ı tutan ve state key'i/attr'ı yeni kalıcı durumu tutan bir nesne döndürMELİdir. Çıplak bir değer veya bir tuple döndürmek çalışmaz — runtime rets["output"] / rets.output ve rets["state"] / rets.state'i arar ve state eksikse fırlatır.

Birden çok çıktılı durum bilgili komponentler için, output bildirilen output stream sayısıyla eşleşen bir tuple'dır. Durum bilgili virtual-output komponentleri için virtual_output ekleyin; streaming modu için pull_request ekleyin.

Platform durumu container yeniden başlatmaları boyunca korur (düzgün crash recovery). İlk başlangıçta, durum initial_state'e geçirilen değerdir.

Saflık ve zamanlama

Komponent gerçekten geçmişe bağlı olmadıkça initial_state kullanmayın. Kalıcı durum olmadan, her çağrı seçilen inputs, mevcut config/files ve deterministik koda göre bağımsızdır. Bu, platformun scheduling/scaling izin verdiğinde güvenle paralelleştirebileceği, batch'leyebileceği veya ölçekleyebileceği şekildir, çünkü bir sonrakini çalıştırmadan önce hiçbir önceki tick'in danışılması gerekmez.

initial_state ayarlandıktan sonra, vertex açıkça durum bilgilidir. Runtime, döndürülen durumu korumalı ve onu bir sonraki tick'e beslemelidir, bu da sıralamayı sözleşmenin bir parçası yapar. Bunu tracker'lar, akümülatörler, session haritaları, pencereler ve aynı girdinin önceki mesajlara bağlı olarak farklı bir çıktı üretebileceği diğer mantık için kullanın.

Modül globalleri model handle'ları, veritabanı client'ları, derlenmiş regex'ler ve yeniden üretilebilir cache'ler için hâlâ yararlıdır. Çıktı geçmişini değiştiren semantik durum içermemelidir. Durum, komponentin ne anlama geldiğini değiştiriyorsa, onu state'e koyun; kod dış dünyayla konuşuyorsa, bu etkileşimi virtual input/output'un arkasına koyun.

Streaming (ileri düzey) — StreamPullShift

1:1 olmayan girdi/çıktı oranlarını işlemesi gereken komponentler için, run(...)'a initial_pull_request= geçin. Fonksiyon o zaman diziler alır (input stream başına bir tane) ve diziler artı bir pull_request (tekil) döndürür.

Bir pull request tek bir seçenek veya sıralı bir seçenekler listesi olabilir. Her seçenek input stream başına bir StreamPullShift içerir. Runtime seçenekleri sırayla kontrol eder ve komponenti şu anda karşılanabilir olan ilk seçenekle çalıştırır.

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), preshift'i atlar, tam olarak size mesajı process'e gösterir, sonra postshift kadar ilerler. StreamPullShift.from_least(preshift, size), preshift'i atlar, sonra en az size mesaj mevcut olduğunda tetiklenir. initial_pull_request= ilk batch'i tohumlar; ondan sonra, her streaming dönüşü pull_request içermelidir.

Sanal girdiler / çıktılar

Virtual input/output, yan etki sınırıdır. Bir komponentin graph'a dönük fonksiyonu scheduler tarafından görülebilir bir stream fonksiyonu olarak tutarken bir side thread, callback, HTTP/WebSocket dinleyici, device reader, model stream'i veya harici gönderici gerektirdiğinde kullanın.

Gelen yan etkiler virtual_in.push(...) üzerinden gider. process fonksiyonu, pull request virtual input'u seçtiğinde onları bir virtual_input anahtar kelime parametresi aracılığıyla alır. Bu payload'lar komponent ve bridge için dahilidir, katı component.yml stream tipleri değildir; bir virtual kanal birden çok payload şekli taşıyabilir, bu nedenle komponent kodu payload'ları açıkça taglemeli, doğrulamalı, dispatch etmeli ve hatalı olanları reddetmelidir. Giden yan etkiler virtual_output üzerinden gider: process work-item'ları döndürür ve virtual_out.set_on_push(...) ile kaydedilen bir bridge callback'i gerçek IO'yu gerçekleştirir.

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],    ],)

İki sıralama kuralı önemlidir:

  • process virtual_input bildiriyorsa, her pull seçeneğinde önce virtual-input shift, sonra bildirilen graph input başına bir shift vardır.
  • use_virtual_output=True ise, process bir virtual_output listesi döndürMELİdir. Durum bilgili komponentler ayrıca state döndürür; streaming komponentleri ayrıca pull_request döndürür.

Component shim'in sizin için yaptıkları

Şunları yapmak zorunda değilsiniz:

  • Stream'leri açmak veya kapatmak.
  • Hat üzerindeki tipli mesajların (de)serileştirmesini ele almak.
  • Aynı node'daki vertex'ler arasında shared memory yönetmek.
  • Health-check'ler veya düzgün kapanma uygulamak.

Şunları yapMANIz gerekir:

  • Fonksiyonu, platformun retry'da yeniden oynatabileceği kadar deterministik yapmak (state argümanı dışında global değiştirilebilir durum yok).
  • Herhangi bir komponent log çıktısı için print() (veya Python'un logging modülü) kullanmak. print(), kanonik komponent loglama yoludur — runtime stdout / stderr'i yakalar ve onu ppl deployment logs üzerinden gösterir.
  • mutable: true parametrelere güvendiğiniz her durumda process() içinde config.sync()'i çağırmak; aksi halde değerler önceki sync'te dondurulmuş kalır.

İçe aktarma kopya kâğıdı

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, tip sistemi kurucularını (Named, List, Union, Tuple) ve registry yardımcılarını (register_named, get_type, create_type) sunar. Bir ad burada yoksa, katalog için type-api/catalog kontrol edin.

Yaygın hatalar

  • numpy'i doğrudan içe aktarmak sorun değil; onu asla requirements.txt'e koymayın — runtime tabanı numpy'i zaten pinler. Onu tekrar pinlemek ABI çakışmaları üretebilir.
  • Platform record'larını ham dict'ler olarak oluşturmayınImage için her zaman pipelogic.cv.Image(...) üzerinden gidin. AudioFrame, Tensor vb. için de aynısı. Elle yapılmış dict'ler, sınırda buffer-vs-shape uyumsuzluklarını yakalayan doğrulamayı atlar.
  • Komponentten sys.exit() / os._exit() kullanmayın — istisnaların yayılmasına izin verin. Shim onları yakalar ve raporlar; açık çıkışlar crash gibi görünür ve retry semantiğini karıştırır. (os.exit() gerçek bir fonksiyon değildir — sys.exit() SystemExit fırlatır ve os._exit() süreci sonlandırır; ikisi de burada yanlıştır.)
  • mutable: true parametreleri okurken process() içinde config.sync()'i unutmak. Çağrı olmadan, değer önceki sync'te (veya modül yüklemede) dondurulmuş kalır, ppl backend change-parameter backend'de başarılı olduktan sonra bile.

İlgili

Bu sayfa yardımcı oldu mu?