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.
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.
Serializes each row of df as one JSON message and publishes it to the topic. The return value is the number of rows sent.
Configuration
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.