Kafka connector reference
This page documents the connector options, table configuration options, and JSON transformer settings for the managed Apache Kafka connector in Lakeflow Connect.
Connector options
The following options configure the Kafka source for each destination table in the ingestion pipeline. Specify these options under connector_options.kafka_options in your pipeline definition. See Examples for full pipeline examples.
Option | Type | Default | Description |
|---|---|---|---|
| List of strings | — | List of topic names to subscribe to. Mutually exclusive with |
| String | — | Java regular expression matching topic names to subscribe to. Mutually exclusive with |
| String |
| Where to begin reading when no checkpoint exists (first run only). Valid values: |
| Transformer | — | Deserializer configuration for message keys. If not set, the key column is retained as |
| Transformer | — | Deserializer configuration for message values. If not set, the value column is retained as |
Transformer options
Transformers define how binary Kafka message keys and values are deserialized into structured columns. Specify the serialization format and the corresponding format-specific options under key_transformer or value_transformer. You can configure a transformer for the key, the value, or both independently. If no transformer is set, the column is retained as BINARY.
For JSON, you can provide an explicit schema, use schema inference with evolution, or omit json_options entirely to store the value as a VARIANT column.
Option | Applies to | Type | Default | Description |
|---|---|---|---|---|
| All | String | — | Serialization format of the data. Valid values: |
| JSON | String | — | Inline schema in Spark DDL format (for example, |
| JSON | String | — | Path to a |
| JSON | String | — | Schema evolution mode for automatic schema inference. See Schema Evolution Modes. |
| JSON | String | — | Comma-separated |
Avro options
Set these options under avro_options when format: AVRO. Provide the schema inline, from a file, or from a schema registry.
Option | Type | Default | Description |
|---|---|---|---|
| String | — | Inline Avro schema in JSON format. Mutually exclusive with |
| String | — | Path to an |
| Object | — | Resolve the schema from a schema registry at runtime instead of |
| String |
| How to handle records that fail to deserialize. Valid values: |
Protobuf options
Set these options under protobuf_options when format: PROTOBUF. Provide a compiled descriptor set (.desc) file and message name, or resolve the schema from a schema registry.
Option | Type | Default | Description |
|---|---|---|---|
| String | — | Path to a compiled Protobuf descriptor set ( |
| String | — | Fully qualified Protobuf message type name (for example, |
| Object | — | Resolve the schema from a schema registry at runtime instead of |
| Integer | — | Maximum expansion depth for recursive Protobuf fields, which Spark SQL does not natively support. Valid values: |
| String |
| How to handle records that fail to deserialize. Valid values: |
Schema registry options
Set schema_registry under avro_options or protobuf_options to resolve the schema at runtime from a Confluent-compatible schema registry. By default, the pipeline authenticates to the registry using the pipeline's Kafka source connection, which stores the registry URL and API key. To authenticate with a different Unity Catalog connection, set connection_name. See Connection properties.
Option | Type | Default | Description |
|---|---|---|---|
| String | — | Required. The subject to resolve in the Confluent-compatible schema registry. |
| String | — | A Unity Catalog connection used to authenticate to the registry. Defaults to the pipeline's Kafka source connection. Set this when the registry uses different credentials. |
| String | — | Protobuf only. Selects a message when the subject defines more than one Protobuf message. Simple ( |
Table configuration options
The following options are set under table_configuration on a table object, a sibling of connector_options. See Examples for full pipeline examples.
Option | Type | Default | Description |
|---|---|---|---|
| String | — | Name of a struct column added to the destination table that holds the Kafka source metadata for each record. See Source metadata column. The name must not be |
Fanout options
This feature is in Private Preview. To try it, reach out to your Databricks contact.
Fanout options route each record from a single Kafka source to one of many destination tables. Specify these options under fanout_options on a schema object (not a table object) in your pipeline definition. See Route records to multiple tables (fanout) for a full pipeline example and Fanout limitations for the constraints.
Option | Type | Default | Description |
|---|---|---|---|
| String | — | Required. SQL expression evaluated against the raw source record (with the Kafka |
| List of Transformer | — | A transform applied to each routed record after routing, before writing to its destination table. Because it runs after routing, it doesn't affect the value that |
Fanout transform options
Each entry in transforms uses the following options. Only JSON format is supported for fanout transforms.
Option | Type | Default | Description |
|---|---|---|---|
| String |
| Optional. Serialization format of the transform. Only |
| String | — | The column the transform reads from and writes back to (for example, |
The output_column transform option is not applied to fanout transforms. The fanout JSON transform always writes its result back to input_column in place; any output_column value is ignored.
Connection properties
When you create the Unity Catalog Kafka connection in Catalog Explorer, you must specify the following properties depending on the authentication method. See Create a Kafka connection for connection creation steps.
Username and password
Property | Description |
|---|---|
Connection name | A unique name for the connection in Unity Catalog. |
Connection type | Select Kafka. |
Auth type | Select Username and Password. |
Username | The username used to authenticate with the Kafka cluster. |
Password | The password used to authenticate with the Kafka cluster. |
Bootstrap servers | The bootstrap server address of the Kafka cluster (for example, |
Schema registry URL (optional) | The URL of your schema registry. |
Schema registry API key (optional) | The API key for your schema registry. |
Schema registry API secret (optional) | The API secret for your schema registry. |
Service Credential
Property | Description |
|---|---|
Connection name | A unique name for the connection in Unity Catalog. |
Connection type | Select Kafka. |
Auth type | Select Service Credential. |
Service credential | Select an existing Unity Catalog service credential or click Create new service credential. |
Bootstrap servers | The bootstrap server address of the Kafka cluster (for example, |
Schema registry URL (optional) | The URL of your schema registry. |
Schema registry API key (optional) | The API key for your schema registry. |
Schema registry API secret (optional) | The API secret for your schema registry. |
Destination table schema
The Kafka connector writes to streaming tables (append-only). The columns written to the destination table depend on whether transformers are configured.
Without transformers (raw binary)
When no key_transformer or value_transformer is configured, the destination table contains the following columns:
Column | Type | Description |
|---|---|---|
|
| The raw binary content of the Kafka message key. |
|
| The raw binary content of the Kafka message value. |
With a STRING transformer
When format: STRING is set on a transformer, the corresponding column is written as STRING instead of BINARY.
With a JSON transformer
When format: JSON is set on a transformer:
- If no
json_optionsare specified, the column is written asVARIANT. - If
json_options.schemaorjson_options.schema_file_pathis specified, the JSON is parsed into typed columns matching the schema. - If
json_options.schema_evolution_modeis set, schema inference is used and the schema evolves automatically.
With an Avro transformer
When format: AVRO is set on a transformer, the column is parsed into typed columns matching the Avro schema. In PERMISSIVE mode (the default), a _corrupt_record column of type BINARY is also added; it holds the raw bytes of any record that fails to deserialize and is null for records that parse successfully.
With a Protobuf transformer
When format: PROTOBUF is set on a transformer, the column is parsed into typed columns matching the Protobuf message definition. In PERMISSIVE mode (the default), a _corrupt_record column of type BINARY is also added; it holds the raw bytes of any record that fails to deserialize and is null for records that parse successfully.
Source metadata column
Set source_metadata_column under table_configuration to add a struct column of that name to the destination table. The struct contains the following Kafka source metadata fields for each record: topic, partition, offset, timestamp, timestampType, and headers. The column name must not be key or value. See Table configuration options.