Standard plugin ยท import "kafka"

Kafka and compatible brokers

Receive JSON or Schema Registry framed Avro records as typed one-row tables, and send DataFrame rows as JSON messages.

JSON and Avro records

kafka::recv

kafka::recv(brokers, topic, group, schema, options = "") -> DataFrame

Polls for one JSON message and decodes its fields using the schema you provide. The schema is a comma-separated list of name:type pairs, for example ts:timestamp,symbol:str,price:f64,volume:i64.

kafka::recv_avro

kafka::recv_avro(brokers, topic, group, schema, registry_url, options = "") -> DataFrame

Reads one Schema Registry framed Avro message, looks up its writer schema from the registry, and decodes a flat record into the fields and types declared in schema.

Both receive functions return one record as a one-row DataFrame. If no message arrives before the poll timeout, the call reports StreamTimeout; use it from an Ibex stream workflow that handles timeouts.

kafka::send

kafka::send(df: DataFrame, brokers: String, topic: String, options: String = "") -> Int

Serializes each row of df as one JSON message and publishes it to the topic. The return value is the number of rows sent.

Options and a basic consumer

The optional options string contains semicolon-separated key=value pairs. For example, poll_timeout_ms=200;consumer.auto.offset.reset=earliest sets the receive timeout and a consumer property. Producers can set flush_timeout_ms and producer properties such as compression.

import "kafka";

let tick = kafka::recv(
    "localhost:19092", "ticks", "ibex-demo",
    "symbol:str,price:f64,size:i64",
    "poll_timeout_ms=1000;consumer.auto.offset.reset=earliest"
);

For broker setup, Schema Registry, full JSON and Avro examples, and the WebSocket dashboard demo, see the Kafka section of the I/O guide.