Skip to content

Kafka Source Connector ​

stream source

The Kafka source connector consumes streaming records from Apache Kafka topics into the rekuiper stream processing engine.

Configuration Overview ​

Configure the Kafka source in $rekuiper/etc/sources/kafka.yaml:

yaml
default:
  brokers: "127.0.0.1:9091,127.0.0.1:9092"
  groupID: ""
  partition: 0
  maxBytes: 1000000

Verify broker reachability before runtime using the Connectivity Check API.

Configuration Properties ​

Property NameOptionalDescription
brokersFalseComma-separated list of Kafka broker addresses (host:port).
saslAuthTypeTrueSASL authentication mechanism: "none", "plain", or "scram". Default is "none".
saslUserNameTrueSASL username credential.
passwordTrueSASL password credential.
insecureSkipVerifyTrueBoolean. Set to true to skip TLS certificate verification.
certificationPathTruePath to client certificate file for mTLS.
privateKeyPathTruePath to client private key file for mTLS.
rootCaPathTruePath to Root CA certificate file.
certficationRawTrueBase64-encoded client certificate string.
privateKeyRawTrueBase64-encoded client private key string.
rootCARawTrueBase64-encoded Root CA certificate string.
maxBytesTrueMaximum bytes fetched per Kafka message batch. Default is 1000000 (1 MB).
groupIDTrueKafka consumer group identifier.
partitionTrueSpecific partition index consumed by the connector.

Create a Stream Source ​

Define a stream using SQL DDL. Set DATASOURCE to the target Kafka topic:

sql
CREATE STREAM kafka_stream () WITH (
  TYPE = "kafka",
  DATASOURCE = "telemetry_topic",
  FORMAT = "json"
);

For REST API and CLI management procedures, refer to Streams Management with REST API and Streams Management with CLI.

Released under the Apache-2.0 / MIT License.