Meddelandetyper

När du inmatar med Zerobus Ingest SDK:erna istället för gRPC väljer du hur poster kodas på en given ström. Zerobus Ingest stöder tre meddelandeformat (JSON, Protocol Buffers (protobuf) och Apache Arrow), så du kan kompromissa enkelhet, typsäkerhet och genomströmning för din arbetsbelastning. Varje format valideras mot ditt Delta-tabellschema innan datan görs hållbar. Se Schemahantering.

Samma 1 000 poster i tre meddelandeformat: protobuf som en kompakt, typad radkodning, Apache Arrow som en kolumnorienterad batch vars metadata och buffertar fördelas över batchen, och JSON som läsbar text med upprepade fältnamn

Vilket format bör du använda?

Format Passar bäst för Notes
JSON Att komma igång och enkla producenter. Det enklaste alternativet, utan någon schema-definition att kompilera. Bekvämt, men långsammare än binära format för högvolymsarbetsbelastningar.
Protocol Buffers Produktion, radorienterade, högvolymsströmmar. Typsäker, kompakt binärkodning. Rekommenderas för de flesta produktionsarbetsbelastningar. Kräver ett kompilerat schema.
Apache Arrow Kolumn- eller batchorienterade arbetsbelastningar. Skickar Apache Arrow-recordbatcher direkt och undviker därmed rad-för-rad-serialisering. Bäst när din data redan är kolumnär eller när du tar in dem i batcher. Se Använda pilflygning med Zerobus-inmatning.

Utsnitten nedan visar formen på varje format i Python SDK. De antar att du redan har skapat SDK-klienten och känner till din måltabell. För hela uppsättningen (endpoint, tabell, tjänsteprincip) och exempel i varje språk, se Use Zerobus Ingest.

JSON

JSON är det enklaste sättet att börja: du skickar poster som JSON-objekt utan schemadefinition att kompilera. Det är idealiskt för att snabbt komma igång, bygga prototyper och för den som värdesätter bekvämlighet mer än ren genomströmning. För högvolymsproduktionsarbetsbelastningar är ett binärt format (protobuf eller Arrow) mer effektivt.

Skapa en ström för JSON-poster genom att skicka ett tabellnamn till TableProperties utan deskriptor:

table_properties = TableProperties(TABLE_NAME)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

stream.ingest_record_offset({"device_name": "sensor-1", "temp": 22, "humidity": 55})

För hela JSON-genomgången, se Skriv en klient.

Protokollbuffertar

Protobuf erbjuder en typsäker, kompakt binär kodning och är det rekommenderade formatet för de flesta produktions- och radorienterade arbetsbelastningar. Du definierar ett protobuf-schema som passar din Delta-måltabell (se Protobuf-schema), kompilerar det och SDK:t importerar poster en post i taget via gRPC.

Att använda protobuf kräver tre steg: generera ett .proto schema som matchar din tabell, kompilera det till en språkmodul, och sedan mata in poster genom att skicka deskriptorn till TableProperties. Följande exempel använder Python SDK.

1. Generera ett .proto schema från din tabell. Python SDK inkluderar ett generate_proto verktyg som läser din Delta-tabell och skriver ett matchande schema:

python -m zerobus.tools.generate_proto \
    --uc-endpoint "https://<workspace-id>.cloud.databricks.com" \
    --client-id "<client-id>" \
    --client-secret "<client-secret>" \
    --table "main.default.air_quality" \
    --output "record.proto" \
    --proto-msg "AirQuality"

Det genererade schemat använder proto2 syntax, med ett valfritt fält för varje Delta-kolumn:

syntax = "proto2";
message AirQuality {
    optional string device_name = 1;
    optional int32 temp = 2;
    optional int64 humidity = 3;
}

2. Kompiliera schemat till en Python-modul med protobuf-kompilatorn:

pip install "grpcio-tools>=1.60.0,<2.0"
python -m grpc_tools.protoc --python_out=. --proto_path=. record.proto

Detta genererar record_pb2.py.

3. Importera poster genom att skicka den kompilerade deskriptorn till TableProperties (standard för protobuf). SDK:n använder deskriptorn för att serialisera varje post:

import record_pb2

descriptor_bytes = record_pb2.AirQuality.DESCRIPTOR.file.serialized_pb
table_properties = TableProperties(TABLE_NAME, descriptor_bytes)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

record = record_pb2.AirQuality(device_name="sensor-1", temp=22, humidity=55)
stream.ingest_record_offset(record)

Exemplet ovan använder Python SDK. Verktygen varierar beroende på språk: vissa SDK:er levererar ett generate_proto verktyg som genererar en .proto från din tabell, medan andra (som Go och TypeScript) kompilerar en befintlig .proto. För stegen per språk, se protobuf-anteckningarna i varje SDK-flik i Write a client. För verktygskällor och fullständiga exempel, se Zerobus SDK-arkivet.

Apache Arrow

Apache Arrow-insamling skickar Arrow-dataRecordBatch direkt över samma gRPC-anslutning, istället för att först konvertera varje rad till JSON eller Protobuf. Det är det bästa valet när din applikation redan producerar Arrow-data eller när du tar in rader i batchar, särskilt för breda, numeriskt tunga eller analysinriktade scheman där rad-för-rad-serialisering lägger till overhead.

Arrow passar också bra för mycket stora satser. Till skillnad från JSON- och protobuf-batchmetoderna, som är allt-eller-inget och begränsas av storleksgränsen per meddelande, delar Arrow Flight-vägen upp en stor batch i mindre transportmeddelanden som skickas och bekräftas individuellt. Se Arrow Flight-batcher är undantaget och använd Arrow Flight med Zerobus Ingest.

Öppna en Arrow-ström med en pyarrow.Schema och mata in RecordBatch data:

stream = sdk.create_arrow_stream(TABLE_NAME, schema, CLIENT_ID, CLIENT_SECRET)

stream.ingest_batch(batch)

För hela Arrow Flight-genomgången, inklusive schemadefinition, batching och komprimering, se Use Arrow Flight with Zerobus Ingest.