// Read from Kafka, ignores the partition and key values
import io.deephaven.kafka.KafkaTools
kafkaProps = new Properties()
kafkaProps.put("bootstrap.servers", "redpanda:9092")
kafkaProps.put("deephaven.partition.column.name", "")
result = KafkaTools.consumeToTable(
kafkaProps,
"share.price",
KafkaTools.ALL_PARTITIONS,
KafkaTools.ALL_PARTITIONS_DONT_SEEK,
KafkaTools.Consume.IGNORE,
KafkaTools.Consume.simpleSpec("Price", double),
KafkaTools.TableType.append()
)
// Read from Kafka, JSON with mapping
import io.deephaven.engine.table.ColumnDefinition
colDefs = [
ColumnDefinition.ofString("Symbol"),
ColumnDefinition.ofString("Side"),
ColumnDefinition.ofDouble("Price"),
ColumnDefinition.ofLong("Qty")
] as ColumnDefinition[]
mapping = [
"jsymbol": "Symbol",
"jside": "Side",
"jprice": "Price",
"jqty": "Qty"
]
result = KafkaTools.consumeToTable(
kafkaProps,
"orders",
KafkaTools.ALL_PARTITIONS,
KafkaTools.ALL_PARTITIONS_DONT_SEEK,
KafkaTools.Consume.IGNORE,
KafkaTools.Consume.jsonSpec(colDefs, mapping, null),
KafkaTools.TableType.append()
)
// Read from Kafka, AVRO
kafkaPropsAvro = new Properties()
kafkaPropsAvro.put("bootstrap.servers", "redpanda:9092")
kafkaPropsAvro.put("schema.registry.url", "http://redpanda:8081")
result = KafkaTools.consumeToTable(
kafkaPropsAvro,
"share.price",
KafkaTools.ALL_PARTITIONS,
KafkaTools.ALL_PARTITIONS_DONT_SEEK,
KafkaTools.Consume.IGNORE,
KafkaTools.Consume.avroSpec("share.price.record"),
KafkaTools.TableType.append()
)