Skip to main content

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

topics

List of strings

—

List of topic names to subscribe to. Mutually exclusive with topic_pattern. Either topics or topic_pattern is required.

topic_pattern

String

—

Java regular expression matching topic names to subscribe to. Mutually exclusive with topics.

starting_offset

String

latest

Where to begin reading when no checkpoint exists (first run only). Valid values: latest, earliest.

key_transformer

Transformer

—

Deserializer configuration for message keys. If not set, the key column is retained as BINARY. See Transformer options.

value_transformer

Transformer

—

Deserializer configuration for message values. If not set, the value column is retained as BINARY. See Transformer options.

Option

Type

Default

Description

topics

List of strings

—

List of topic names to subscribe to. Mutually exclusive with topic_pattern. Either topics or topic_pattern is required.

topic_pattern

String

—

Java regular expression matching topic names to subscribe to. Mutually exclusive with topics.

starting_offset

String

latest

Where to begin reading when no checkpoint exists (first run only). Valid values: latest, earliest.

key_transformer

Transformer

—

Deserializer configuration for message keys. If not set, the key column is retained as BINARY. See Transformer options.

value_transformer

Transformer

—

Deserializer configuration for message values. If not set, the value column is retained as BINARY. See Transformer options.

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

format

All

String

—

Serialization format of the data. Valid values: STRING, JSON, AVRO, PROTOBUF. STRING requires no additional options. If no json_options are specified on the transformer, the value is parsed as VARIANT by default. See Variant Data Format for more information. For AVRO and PROTOBUF, see Avro options and Protobuf options.

json_options.schema

JSON

String

—

Inline schema in Spark DDL format (for example, "id BIGINT, name STRING"). Mutually exclusive with schema_file_path.

json_options.schema_file_path

JSON

String

—

Path to a .ddl schema file. Mutually exclusive with schema. Supports Unity Catalog Volumes paths (/Volumes/...).

json_options.schema_evolution_mode

JSON

String

—

Schema evolution mode for automatic schema inference. See Schema Evolution Modes.

json_options.schema_hints

JSON

String

—

Comma-separated "column_name type" pairs to influence schema inference (for example, "id BIGINT, ts TIMESTAMP"). Requires schema_evolution_mode to be set. See Override schema inference using schema hints.

Option

Applies to

Type

Default

Description

format

All

String

—

Serialization format of the data. Valid values: STRING, JSON, AVRO, PROTOBUF. STRING requires no additional options. If no json_options are specified on the transformer, the value is parsed as VARIANT by default. See Variant Data Format for more information. For AVRO and PROTOBUF, see Avro options and Protobuf options.

json_options.schema

JSON

String

—

Inline schema in Spark DDL format (for example, "id BIGINT, name STRING"). Mutually exclusive with schema_file_path.

json_options.schema_file_path

JSON

String

—

Path to a .ddl schema file. Mutually exclusive with schema. Supports Unity Catalog Volumes paths (/Volumes/...).

json_options.schema_evolution_mode

JSON

String

—

Schema evolution mode for automatic schema inference. See Schema Evolution Modes.

json_options.schema_hints

JSON

String

—

Comma-separated "column_name type" pairs to influence schema inference (for example, "id BIGINT, ts TIMESTAMP"). Requires schema_evolution_mode to be set. See Override schema inference using schema hints.

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

avro_options.schema

String

—

Inline Avro schema in JSON format. Mutually exclusive with schema_file_path and schema_registry.

avro_options.schema_file_path

String

—

Path to an .avsc schema file. Supports Unity Catalog Volumes paths (/Volumes/...). Mutually exclusive with schema and schema_registry.

avro_options.schema_registry

Object

—

Resolve the schema from a schema registry at runtime instead of schema or schema_file_path. See Schema registry options.

avro_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valid values: PERMISSIVE (default), which writes the raw bytes of each record that fails to deserialize to a _corrupt_record column (type BINARY), sets the record's other columns to null, and continues processing; and FAILFAST, which fails the pipeline on the first record that fails to deserialize.

Option

Type

Default

Description

avro_options.schema

String

—

Inline Avro schema in JSON format. Mutually exclusive with schema_file_path and schema_registry.

avro_options.schema_file_path

String

—

Path to an .avsc schema file. Supports Unity Catalog Volumes paths (/Volumes/...). Mutually exclusive with schema and schema_registry.

avro_options.schema_registry

Object

—

Resolve the schema from a schema registry at runtime instead of schema or schema_file_path. See Schema registry options.

avro_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valid values: PERMISSIVE (default), which writes the raw bytes of each record that fails to deserialize to a _corrupt_record column (type BINARY), sets the record's other columns to null, and continues processing; and FAILFAST, which fails the pipeline on the first record that fails to deserialize.

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

protobuf_options.desc_file_path

String

—

Path to a compiled Protobuf descriptor set (.desc) file. Supports Unity Catalog Volumes paths (/Volumes/...). Required unless schema_registry is set.

protobuf_options.message_name

String

—

Fully qualified Protobuf message type name (for example, com.example.events.UserEvent). Required with desc_file_path.

protobuf_options.schema_registry

Object

—

Resolve the schema from a schema registry at runtime instead of desc_file_path. See Schema registry options.

protobuf_options.recursive_fields_max_depth

Integer

—

Maximum expansion depth for recursive Protobuf fields, which Spark SQL does not natively support. Valid values: -1 (disallow recursive fields), 0 (drop them), 1–10.

protobuf_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valid values: PERMISSIVE (default), which writes the raw bytes of each record that fails to deserialize to a _corrupt_record column (type BINARY), sets the record's other columns to null, and continues processing; and FAILFAST, which fails the pipeline on the first record that fails to deserialize.

Option

Type

Default

Description

protobuf_options.desc_file_path

String

—

Path to a compiled Protobuf descriptor set (.desc) file. Supports Unity Catalog Volumes paths (/Volumes/...). Required unless schema_registry is set.

protobuf_options.message_name

String

—

Fully qualified Protobuf message type name (for example, com.example.events.UserEvent). Required with desc_file_path.

protobuf_options.schema_registry

Object

—

Resolve the schema from a schema registry at runtime instead of desc_file_path. See Schema registry options.

protobuf_options.recursive_fields_max_depth

Integer

—

Maximum expansion depth for recursive Protobuf fields, which Spark SQL does not natively support. Valid values: -1 (disallow recursive fields), 0 (drop them), 1–10.

protobuf_options.parse_mode

String

PERMISSIVE

How to handle records that fail to deserialize. Valid values: PERMISSIVE (default), which writes the raw bytes of each record that fails to deserialize to a _corrupt_record column (type BINARY), sets the record's other columns to null, and continues processing; and FAILFAST, which fails the pipeline on the first record that fails to deserialize.

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

schema_registry.confluent_options.subject

String

—

Required. The subject to resolve in the Confluent-compatible schema registry.

schema_registry.connection_name

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.

schema_registry.protobuf_message_name

String

—

Protobuf only. Selects a message when the subject defines more than one Protobuf message. Simple (Location) or fully qualified (com.example.protos.Location). Defaults to the first message in the schema.

Option

Type

Default

Description

schema_registry.confluent_options.subject

String

—

Required. The subject to resolve in the Confluent-compatible schema registry.

schema_registry.connection_name

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.

schema_registry.protobuf_message_name

String

—

Protobuf only. Selects a message when the subject defines more than one Protobuf message. Simple (Location) or fully qualified (com.example.protos.Location). Defaults to the first message in the schema.

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

source_metadata_column

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 key or value.

Option

Type

Default

Description

source_metadata_column

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 key or value.

Fanout options​

Preview

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

fanout_by

String

—

Required. SQL expression evaluated against the raw source record (with the Kafka key and value columns) whose result determines the destination table name. The value becomes the final segment of the table name: {destination_catalog}.{destination_schema}.{value}. The value column is binary, so cast it to a string before extracting a field, for example cast(value as string):event_type::string. The expression must resolve to a non-null STRING, and that string is used as the table-name segment verbatim, without quoting or sanitization, so it must be a valid unquoted table identifier. Values with spaces, dots, or other characters that aren't valid in an unquoted identifier fail the write, as does a value that is entirely digits (for example, 123). A leading digit is allowed when the value also contains a letter or underscore, such as 2024_events. Destination tables are created automatically if they don't already exist.

transforms

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_by sees. At most one transform is allowed, and it must use format: JSON. See Fanout transform options.

Option

Type

Default

Description

fanout_by

String

—

Required. SQL expression evaluated against the raw source record (with the Kafka key and value columns) whose result determines the destination table name. The value becomes the final segment of the table name: {destination_catalog}.{destination_schema}.{value}. The value column is binary, so cast it to a string before extracting a field, for example cast(value as string):event_type::string. The expression must resolve to a non-null STRING, and that string is used as the table-name segment verbatim, without quoting or sanitization, so it must be a valid unquoted table identifier. Values with spaces, dots, or other characters that aren't valid in an unquoted identifier fail the write, as does a value that is entirely digits (for example, 123). A leading digit is allowed when the value also contains a letter or underscore, such as 2024_events. Destination tables are created automatically if they don't already exist.

transforms

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_by sees. At most one transform is allowed, and it must use format: JSON. See Fanout transform options.

Fanout transform options​

Each entry in transforms uses the following options. Only JSON format is supported for fanout transforms.

Option

Type

Default

Description

format

String

JSON

Optional. Serialization format of the transform. Only JSON is supported for fanout transforms, and it is the default when omitted. STRING, AVRO, and PROTOBUF are not supported.

input_column

String

—

The column the transform reads from and writes back to (for example, value). The fanout JSON transform parses this column in place.

Option

Type

Default

Description

format

String

JSON

Optional. Serialization format of the transform. Only JSON is supported for fanout transforms, and it is the default when omitted. STRING, AVRO, and PROTOBUF are not supported.

input_column

String

—

The column the transform reads from and writes back to (for example, value). The fanout JSON transform parses this column in place.

note

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, broker1:9092,broker2:9092).

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.

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, broker1:9092,broker2:9092).

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, broker1:9092,broker2:9092).

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.

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, broker1:9092,broker2:9092).

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

key

BINARY

The raw binary content of the Kafka message key.

value

BINARY

The raw binary content of the Kafka message value.

Column

Type

Description

key

BINARY

The raw binary content of the Kafka message key.

value

BINARY

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_options are specified, the column is written as VARIANT.
  • If json_options.schema or json_options.schema_file_path is specified, the JSON is parsed into typed columns matching the schema.
  • If json_options.schema_evolution_mode is 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.