Skip to main content

Kafka

Overview

Apache Kafka is a distributed event streaming platform used for high-throughput, fault-tolerant messaging and real-time data pipelines. Confluent Cloud provides a fully managed Kafka service with an integrated Schema Registry.

The DataHub Kafka connector extracts topic metadata and schemas from Kafka clusters and Confluent Schema Registry. It supports Avro, Protobuf, and JSON schemas, multi-stage schema resolution with automatic fallback and inference for topics not registered in the schema registry, and optional data profiling to generate field-level statistics and sample values from message content. The integration also captures stateful deletion detection.

Concept Mapping

Kafka ConceptDataHub ConceptNotes
TopicDatasetSubtype Topic
Schema (Subject)Schema MetadataAvro, Protobuf, JSON schemas
Message FieldsDataset FieldsExtracted from schemas or inferred
Kafka ClusterData Platform InstanceWhen platform_instance is configured
Schema metadataTags, Glossary Terms, Owners (CorpUser / CorpGroup)Optional, Avro only: derived from schema properties when enable_meta_mapping is set with meta_mapping / field_meta_mapping directives

Module kafka

Certified

Important Capabilities

CapabilityStatusNotes
Column-level LineageNot supported.
Data ProfilingOptionally enabled via configuration profiling.enabled.
DescriptionsSet dataset description to top level doc field for Avro schema.
Detect Deleted EntitiesEnabled by default via stateful ingestion.
Platform InstanceFor multiple Kafka clusters, use the platform_instance configuration.
Schema MetadataSchemas associated with each topic are extracted from the schema registry. Avro and Protobuf (certified), JSON (incubating). Schema references are supported.
Table-Level LineageNot supported. If you use Kafka Connect, the kafka-connect source can generate lineage.
Test ConnectionEnabled by default.

Overview

The kafka module ingests metadata from Kafka into DataHub. It is intended for production ingestion workflows and module-specific capabilities are documented below.

Extract Topics & Schemas from Apache Kafka or Confluent Cloud.

This plugin extracts the following:

  • Topics from the Kafka broker
  • Schemas associated with each topic from the schema registry (Avro, Protobuf and JSON schemas are supported)

Prerequisites

Before running ingestion, ensure network connectivity to the source, valid authentication credentials, and read permissions for metadata APIs required by this module.

Install the Plugin

pip install 'acryl-datahub[kafka]'

Starter Recipe

Check out the following recipe to get started with ingestion! See below for full configuration options.

For general pointers on writing and running a recipe, see our main recipe guide.

source:
type: "kafka"
config:
platform_instance: "YOUR_CLUSTER_ID"
connection:
bootstrap: "broker:9092"
schema_registry_url: http://localhost:8081

# Optional: Enable data profiling
profiling:
enabled: true
sample_size: 1000
max_sample_time_seconds: 60
sampling_strategy: "latest"

sink:
# sink configs

Config Details

Note that a . is used to denote nested fields in the YAML recipe.

FieldDescription
convert_urns_to_lowercase
boolean
Whether to convert dataset urns to lowercase. This value is part of each dataset's URN identity, so it must stay fixed for the life of a deployment. Changing it after data has been ingested re-keys every dataset (e.g. MyDb.MyTable becomes mydb.mytable); with stateful ingestion enabled the old-cased URNs are then soft-deleted as stale while the new-cased ones are created, producing duplicate or orphaned entities. Pick one value before the first run and leave it unchanged.
Default: False
disable_topic_record_naming_strategy
boolean
Disables the utilization of the TopicRecordNameStrategy for Schema Registry subjects. For more information, visit: https://docs.confluent.io/platform/current/schema-registry/serdes-develop/index.html#handling-differences-between-preregistered-and-client-derived-schemas:~:text=io.confluent.kafka.serializers.subject.TopicRecordNameStrategy
Default: False
enable_meta_mapping
boolean
When enabled, applies the mappings that are defined through the meta_mapping directives.
Default: True
external_url_base
One of string, null
Base URL for external platform (e.g. Aiven) where topics can be viewed. The topic name will be appended to this base URL.
Default: None
field_meta_mapping
object
mapping rules that will be executed against field-level schema properties. Refer to the section below on meta automated mappings.
Default: {}
ignore_warnings_on_schema_type
boolean
Disables warnings reported for non-AVRO/Protobuf value or key schemas if set.
Default: False
ingest_schemas_as_entities
boolean
Enables ingesting schemas from schema registry as separate entities, in addition to the topics
Default: False
meta_mapping
object
mapping rules that will be executed against top-level schema properties. Refer to the section below on meta automated mappings.
Default: {}
platform_instance
One of string, null
The instance of the platform that all assets produced by this recipe belong to. This should be unique within the platform. See https://docs.datahub.com/docs/platform-instances/ for more details.
Default: None
schema_registry_class
string
The fully qualified implementation class(custom) that implements the KafkaSchemaRegistryBase interface.
Default: datahub.ingestion.source.confluent_schema_registry...
schema_tags_field
string
The field name in the schema metadata that contains the tags to be added to the dataset.
Default: tags
strip_user_ids_from_email
boolean
Whether or not to strip email id while adding owners using meta mappings.
Default: False
tag_prefix
string
Prefix added to tags during ingestion.
Default:
topic_subject_map
map(str,string)
env
string
The environment that all assets produced by this connector belong to
Default: PROD
connection
KafkaConsumerConnectionConfig
Configuration class for holding connectivity information for Kafka consumers
connection.bootstrap
string
Default: localhost:9092
connection.client_timeout_seconds
integer
The request timeout used when interacting with the Kafka APIs.
Default: 60
connection.consumer_config
object
Extra consumer config serialized as JSON. These options will be passed into Kafka's DeserializingConsumer. See https://docs.confluent.io/platform/current/clients/confluent-kafka-python/html/index.html#deserializingconsumer and https://github.com/edenhill/librdkafka/blob/master/CONFIGURATION.md .
connection.schema_registry_config
object
Extra schema registry config serialized as JSON. These options will be passed into Kafka's SchemaRegistryClient. https://docs.confluent.io/platform/current/clients/confluent-kafka-python/html/index.html?#schemaregistryclient
connection.schema_registry_url
string
Schema registry URL. Can be overridden with KAFKA_SCHEMAREGISTRY_URL environment variable, or will use DATAHUB_GMS_BASE_PATH if not set.
domain
map(str,AllowDenyPattern)
A class to store allow deny regexes
domain.key.allow
array
List of regex patterns to include in ingestion
Default: ['.*']
domain.key.allow.string
string
domain.key.ignoreCase
One of boolean, null
Whether to ignore case sensitivity during pattern matching.
Default: True
domain.key.deny
array
List of regex patterns to exclude from ingestion.
Default: []
domain.key.deny.string
string
schema_resolution
SchemaResolutionFallback
schema_resolution.enabled
boolean
Enable comprehensive schema resolution with multiple fallback strategies for topics where schema registry lookup fails.
Default: False
schema_resolution.max_messages_per_topic
integer
Maximum number of messages to sample per topic for record name extraction and schema inference. Must be positive.
Default: 10
schema_resolution.offset_reset_strategy
Enum
One of: "earliest", "latest", "hybrid"
schema_resolution.sample_timeout_seconds
number
Maximum time to spend sampling messages from a single topic (in seconds) for record name extraction and schema inference. Must be positive.
Default: 2.0
topic_patterns
AllowDenyPattern
A class to store allow deny regexes
topic_patterns.ignoreCase
One of boolean, null
Whether to ignore case sensitivity during pattern matching.
Default: True
topic_patterns.allow
array
List of regex patterns to include in ingestion
Default: ['.*']
topic_patterns.allow.string
string
topic_patterns.deny
array
List of regex patterns to exclude from ingestion.
Default: []
topic_patterns.deny.string
string
profiling
ProfilerConfig
profiling.batch_size
integer
Number of messages to fetch in a single batch (for more efficient reading). Must be positive.
Default: 100
profiling.catch_exceptions
boolean
Default: True
profiling.enabled
boolean
Whether profiling should be done.
Default: False
profiling.field_sample_values_limit
integer
Upper limit for number of sample values to collect for all columns.
Default: 20
profiling.include_field_distinct_count
boolean
Whether to profile for the number of distinct values for each column.
Default: True
profiling.include_field_distinct_value_frequencies
boolean
Whether to profile for distinct value frequencies.
Default: False
profiling.include_field_histogram
boolean
Whether to profile for the histogram for numeric fields.
Default: False
profiling.include_field_max_value
boolean
Whether to profile for the max value of numeric columns.
Default: True
profiling.include_field_mean_value
boolean
Whether to profile for the mean value of numeric columns.
Default: True
profiling.include_field_median_value
boolean
Whether to profile for the median value of numeric columns.
Default: True
profiling.include_field_min_value
boolean
Whether to profile for the min value of numeric columns.
Default: True
profiling.include_field_null_count
boolean
Whether to profile for the number of nulls for each column.
Default: True
profiling.include_field_quantiles
boolean
Whether to profile for the quantiles of numeric columns.
Default: False
profiling.include_field_sample_values
boolean
Whether to profile for the sample values for all columns.
Default: True
profiling.include_field_stddev_value
boolean
Whether to profile for the standard deviation of numeric columns.
Default: True
profiling.limit
One of integer, null
Max number of documents to profile. By default, profiles all documents.
Default: None
profiling.max_number_of_fields_to_profile
One of integer, null
A positive integer that specifies the maximum number of columns to profile for any table. None implies all columns. The cost of profiling goes up significantly as the number of columns to profile goes up.
Default: None
profiling.max_sample_time_seconds
integer
Maximum time to spend sampling messages in seconds. Must be positive.
Default: 60
profiling.max_workers
integer
Number of worker threads to use for profiling. Set to 1 to disable.
Default: 20
profiling.method
Enum
One of: "ge", "sqlalchemy"
Default: sqlalchemy
profiling.nested_field_max_depth
integer
Maximum recursion depth when flattening nested JSON structures during profiling. Lower values prevent recursion errors but may truncate deeply nested data. Applies to connectors that process dynamic JSON content (e.g., Kafka, MongoDB, Elasticsearch).
Default: 10
profiling.offset
One of integer, null
Offset in documents to profile. By default, uses no offset.
Default: None
profiling.partition_datetime
One of string(date-time), null
If specified, profile only the partition which matches this datetime. If not specified, profile the latest partition. Only Bigquery supports this.
Default: None
profiling.partition_profiling_enabled
boolean
Whether to profile partitioned tables. Only BigQuery and Aws Athena supports this. If enabled, latest partition data is used for profiling.
Default: True
profiling.profile_external_tables
boolean
Whether to profile external tables. Only Snowflake and Redshift supports this.
Default: False
profiling.profile_if_updated_since_days
One of number, null
Profile table only if it has been updated since these many number of days. If set to null, no constraint of last modified time for tables to profile. Supported in Snowflake, BigQuery, and Dremio. Note: for Dremio this compares against DataHub's last-profiled timestamp (Dremio exposes no table modification time), so it controls profile frequency rather than reacting to upstream change.
Default: None
profiling.profile_nested_fields
boolean
Whether to profile complex types like structs, arrays and maps.
Default: False
profiling.profile_table_level_only
boolean
Whether to perform profiling at table-level only, or include column-level profiling as well.
Default: False
profiling.profile_table_row_count_estimate_only
boolean
Use an approximate query for row count. This will be much faster but slightly less accurate. Only supported for Postgres and MySQL.
Default: False
profiling.profile_table_row_limit
One of integer, null
Profile tables only if their row count is less than specified count. If set to null, no limit on the row count of tables to profile. Supported only in Snowflake, BigQuery. Supported for Oracle based on gathered stats.
Default: 5000000
profiling.profile_table_size_limit
One of integer, null
Profile tables only if their size is less than specified GBs. If set to null, no limit on the size of tables to profile. Supported in Snowflake, BigQuery, Databricks, Oracle, and Teradata. Oracle uses calculated size from gathered stats. Teradata uses DBC space accounting.
Default: 5
profiling.query_combiner_enabled
boolean
This feature is still experimental and can be disabled if it causes issues. Reduces the total number of queries issued and speeds up profiling by dynamically combining SQL queries where possible.
Default: True
profiling.report_dropped_profiles
boolean
Whether to report datasets or dataset columns which were not profiled. Set to True for debugging purposes.
Default: False
profiling.sample_size
integer
Number of messages to sample for profiling. Higher values provide more accurate statistics but take longer to process. Must be positive.
Default: 200
profiling.sampling_strategy
Enum
One of: "latest", "random", "stratified", "full"
profiling.turn_off_expensive_profiling_metrics
boolean
Whether to turn off expensive profiling or not. This turns off profiling for quantiles, distinct_value_frequencies, histogram & sample_values. This also limits maximum number of fields being profiled to 10.
Default: False
profiling.use_sampling
boolean
Whether to profile column level stats on sample of table. Only BigQuery and Snowflake support this. If enabled, profiling is done on rows sampled from table. Sampling is not done for smaller tables.
Default: True
profiling.operation_config
OperationConfig
profiling.operation_config.lower_freq_profile_enabled
boolean
Whether to do profiling at lower freq or not. This does not do any scheduling just adds additional checks to when not to run profiling.
Default: False
profiling.operation_config.profile_date_of_month
One of integer, null
Number between 1 to 31 for date of month (both inclusive). If not specified, defaults to Nothing and this field does not take affect.
Default: None
profiling.operation_config.profile_day_of_week
One of integer, null
Number between 0 to 6 for day of week (both inclusive). 0 is Monday and 6 is Sunday. If not specified, defaults to Nothing and this field does not take affect.
Default: None
profiling.tags_to_ignore_sampling
One of array, null
Fixed list of tags to ignore sampling. Each entry may be a full tag URN (e.g. urn:li:tag:my_tag) or just the tag name (e.g. my_tag). If not specified, tables will be sampled based on use_sampling.
Default: None
profiling.tags_to_ignore_sampling.string
string
stateful_ingestion
One of StatefulStaleMetadataRemovalConfig, null
Default: None
stateful_ingestion.enabled
boolean
Whether or not to enable stateful ingest. Default: True if a pipeline_name is set and either a datahub-rest sink or datahub_api is specified, otherwise False
Default: False
stateful_ingestion.fail_safe_threshold
number
Prevents large amount of soft deletes & the state from committing from accidental changes to the source configuration if the relative change percent in entities compared to the previous state is above the 'fail_safe_threshold'.
Default: 75.0
stateful_ingestion.remove_stale_metadata
boolean
Soft-deletes the entities present in the last successful run but missing in the current run with stateful_ingestion enabled.
Default: True

Capabilities

note

Stateful Ingestion is available only when a Platform Instance is assigned to this source.

Use the Important Capabilities table above as the source of truth for supported features and whether additional configuration is required.

Connecting to Confluent Cloud

If using Confluent Cloud you can use a recipe like this. In this consumer_config.sasl.username and consumer_config.sasl.password are the API credentials that you get (in the Confluent UI) from your cluster -> Data Integration -> API Keys. schema_registry_config.basic.auth.user.info has API credentials for Confluent schema registry which you get (in Confluent UI) from Schema Registry -> API credentials.

When creating API Key for the cluster ensure that the ACLs associated with the key are set like below. This is required for DataHub to read topic metadata from topics in Confluent Cloud.

Topic Name = *
Permission = ALLOW
Operation = DESCRIBE
Pattern Type = LITERAL
source:
type: "kafka"
config:
platform_instance: "YOUR_CLUSTER_ID"
connection:
bootstrap: "abc-defg.eu-west-1.aws.confluent.cloud:9092"
consumer_config:
security.protocol: "SASL_SSL"
sasl.mechanism: "PLAIN"
sasl.username: "${CLUSTER_API_KEY_ID}"
sasl.password: "${CLUSTER_API_KEY_SECRET}"
schema_registry_url: "https://abc-defgh.us-east-2.aws.confluent.cloud"
schema_registry_config:
basic.auth.user.info: "${REGISTRY_API_KEY_ID}:${REGISTRY_API_KEY_SECRET}"

sink:
# sink configs

If you are trying to add domains to your topics you can use a configuration like below.

source:
type: "kafka"
config:
# ...connection block
domain:
"urn:li:domain:13ae4d85-d955-49fc-8474-9004c663a810":
allow:
- ".*"
"urn:li:domain:d6ec9868-6736-4b1f-8aa6-fee4c5948f17":
deny:
- ".*"

Note that the domain in config above can be either an urn or a domain id (i.e. urn:li:domain:13ae4d85-d955-49fc-8474-9004c663a810 or simply 13ae4d85-d955-49fc-8474-9004c663a810). The Domain should exist in your DataHub instance before ingesting data into the Domain. To create a Domain on DataHub, check out the Domains User Guide.

If you are using a non-default subject naming strategy in the schema registry, such as RecordNameStrategy, the mapping for the topic's key and value schemas to the schema registry subject names should be provided via topic_subject_map as shown in the configuration below.

source:
type: "kafka"
config:
# ...connection block
# Defines the mapping for the key & value schemas associated with a topic & the subject name registered with the
# kafka schema registry.
topic_subject_map:
# Defines both key & value schema for topic 'my_topic_1'
"my_topic_1-key": "io.acryl.Schema1"
"my_topic_1-value": "io.acryl.Schema2"
# Defines only the value schema for topic 'my_topic_2' (the topic doesn't have a key schema).
"my_topic_2-value": "io.acryl.Schema3"

Custom Schema Registry

The Kafka Source uses the schema registry to figure out the schema associated with both key and value for the topic. By default it uses the Confluent's Kafka Schema registry and supports the AVRO and PROTOBUF schema types.

If you're using a custom schema registry, or you are using schema type other than AVRO or PROTOBUF, then you can provide your own custom implementation of the KafkaSchemaRegistryBase class, and implement the get_schema_metadata(topic, platform_urn) method that given a topic name would return object of SchemaMetadata containing schema for that topic. Please refer datahub.ingestion.source.confluent_schema_registry::ConfluentSchemaRegistry for sample implementation of this class.

class KafkaSchemaRegistryBase(ABC):
@abstractmethod
def get_schema_metadata(
self, topic: str, platform_urn: str
) -> Optional[SchemaMetadata]:
pass

The custom schema registry class can be configured using the schema_registry_class config param of the kafka source as shown below.

source:
type: "kafka"
config:
# Set the custom schema registry implementation class
schema_registry_class: "datahub.ingestion.source.confluent_schema_registry.ConfluentSchemaRegistry"
# Coordinates
connection:
bootstrap: "broker:9092"
schema_registry_url: http://localhost:8081

OAuth Callback

The OAuth callback function can be set up for both Kafka sources (consumers) and sinks (producers):

  • For sources: config.connection.consumer_config.oauth_cb
  • For sinks: config.connection.producer_config.oauth_cb

You need to specify a Python function reference in the format <python-module>:<function-name>.

For example, in the configuration oauth:create_token, create_token is a function defined in oauth.py, and oauth.py must be accessible in the PYTHONPATH.

Deploying Custom OAuth Callbacks

For Built-in Callbacks (Recommended):

DataHub includes pre-built OAuth callbacks for common use cases:

  • AWS MSK IAM: datahub_actions.utils.kafka_msk_iam:oauth_cb
  • Azure Event Hubs: datahub_actions.utils.kafka_eventhubs_auth:oauth_cb

Important: To use these built-in callbacks, you must install the acryl-datahub-actions package:

pip install acryl-datahub-actions>=1.3.1.2

For Custom OAuth Callbacks:

If you need to implement a custom OAuth callback, you must ensure your Python module is accessible to the DataHub process, e.g. adding it via PYTHONPATH=/path/to/your/module:$PYTHONPATH or pip install my-oauth-package.

Example for Kafka Source:

source:
type: "kafka"
config:
# Set the custom schema registry implementation class
schema_registry_class: "datahub.ingestion.source.confluent_schema_registry.ConfluentSchemaRegistry"
# Coordinates
connection:
bootstrap: "broker:9092"
schema_registry_url: http://localhost:8081
consumer_config:
security.protocol: "SASL_PLAINTEXT"
sasl.mechanism: "OAUTHBEARER"
oauth_cb: "oauth:create_token"
# sink configs

Example for Kafka Sink (e.g., MSK IAM authentication):

sink:
type: "datahub-kafka"
config:
connection:
bootstrap: "b-1.msk.us-west-2.amazonaws.com:9098"
schema_registry_url: "http://datahub-gms:8080/schema-registry/api/"
producer_config:
security.protocol: "SASL_SSL"
sasl.mechanism: "OAUTHBEARER"
sasl.oauthbearer.method: "default"
oauth_cb: "datahub_actions.utils.kafka_msk_iam:oauth_cb"

Enriching DataHub metadata with automated meta mapping

note

Meta mapping is currently only available for Avro schemas, and requires that those Avro schemas are pushed to the schema registry.

Avro schemas are permitted to have additional attributes not defined by the specification as arbitrary metadata. A common pattern is to utilize this for business metadata. The Kafka source has the ability to transform this directly into DataHub Owners, Tags and Terms.

Simple tags

If you simply have a list of tags embedded into an Avro schema (either at the top-level or for an individual field), you can use the schema_tags_field config.

Example Avro schema:

{
"name": "sampleRecord",
"type": "record",
"tags": ["tag1", "tag2"],
"fields": [
{
"name": "field_1",
"type": "string",
"tags": ["tag3", "tag4"]
}
]
}

The name of the field containing a list of tags can be configured with the schema_tags_field property:

config:
schema_tags_field: tags
meta mapping

You can also map specific Avro fields into Owners, Tags and Terms using meta mapping.

Example Avro schema:

{
"name": "sampleRecord",
"type": "record",
"owning_team": "@Data-Science",
"data_tier": "Bronze",
"fields": [
{
"name": "field_1",
"type": "string",
"gdpr": {
"pii": true
}
}
]
}

This can be mapped to DataHub metadata with meta_mapping config:

config:
meta_mapping:
owning_team:
match: "^@(.*)"
operation: "add_owner"
config:
owner_type: group
data_tier:
match: "Bronze|Silver|Gold"
operation: "add_term"
config:
term: "{{ $match }}"
field_meta_mapping:
gdpr.pii:
match: true
operation: "add_tag"
config:
tag: "pii"

The underlying implementation is similar to dbt meta mapping, which has more detailed examples that can be used for reference.

Schema Resolution

DataHub provides multi-stage schema resolution for topics whose schemas are not registered in the schema registry, or that use non-default naming strategies. This feature is independent of data profiling.

When schema_resolution.enabled is true, DataHub attempts to resolve schemas using the following stages in order:

  1. TopicNameStrategy: Direct lookup using <topic>-key/value pattern (most common)
  2. TopicSubjectMap: User-defined topic-to-subject mappings via topic_subject_map config
  3. RecordNameStrategy: Extract record names from message data and lookup <record_name>-key/value
  4. TopicRecordNameStrategy: Combine topic + record name: <topic>-<record_name>-key/value
  5. Schema Inference: As final fallback, infer schema from message data analysis

This ensures maximum compatibility with different Confluent Schema Registry naming strategies.

source:
type: kafka
config:
schema_resolution:
enabled: true # disabled by default
sample_timeout_seconds: 2.0
offset_reset_strategy: "hybrid" # "earliest", "latest", or "hybrid"
max_messages_per_topic: 10

profiling:
max_workers: 20 # controls parallelization for both profiling and schema resolution
nested_field_max_depth: 5

Sampling strategies for schema inference:

  • hybrid (default): Tries latest first for speed, falls back to earliest if no recent messages found.
  • latest: Only reads recent messages. Fastest but may fail on quiet topics.
  • earliest: Scans from the beginning of topic history. Most comprehensive but slower on large topics.

Performance notes:

  • Inferred schemas are cached for 60 minutes to avoid re-sampling.
  • Worker count scales automatically based on CPU cores and topic count.
  • Set schema_resolution.enabled: false to receive warnings for missing schemas instead of automatic resolution.

Data Profiling

The Kafka source supports data profiling of message content to generate field-level statistics and sample values. Profiling and schema resolution are enabled independently — either can be turned on without the other. Note, however, that both features draw their parallelism from the same profiling.max_workers setting.

source:
type: "kafka"
config:
profiling:
enabled: true
sample_size: 200 # messages to sample per topic
max_sample_time_seconds: 60
sampling_strategy: "latest" # latest, random, stratified, or full
max_workers: 4
batch_size: 100

# Field-level statistics
include_field_null_count: true
include_field_distinct_count: true
include_field_min_value: true
include_field_max_value: true
include_field_mean_value: true
include_field_median_value: true
include_field_stddev_value: true
include_field_quantiles: false # expensive, disabled by default
include_field_distinct_value_frequencies: false # expensive
include_field_histogram: false # expensive
include_field_sample_values: true

# Nested field handling
profile_nested_fields: true
nested_field_max_depth: 10

# Scheduled profiling (optional)
operation_config:
lower_freq_profile_enabled: false
profile_day_of_week: 1 # Monday=0, Sunday=6
profile_date_of_month: 15

Sampling strategies:

  • latest (default): Samples the most recent messages from the end of each partition.
  • random: Samples messages from random offsets across partitions.
  • stratified: Evenly distributes samples across the topic timeline.
  • full: Processes the entire topic (respects sample_size limit).

The nested_field_max_depth setting (default: 10) prevents recursion errors on deeply nested or circular JSON structures. Reduce it for topics with complex nested messages.

Limitations

Module behavior is constrained by source APIs, permissions, and metadata exposed by the platform. Refer to capability notes for unsupported or conditional features.

PROTOBUF Schema Type Limitations

The current implementation of the support for PROTOBUF schema type has the following limitations:

  • Recursive types are not supported.
  • Protobuf compilation uses a process-global descriptor pool. The first topic that compiles a given fully-qualified message type succeeds; any later topic that embeds the same type name hits a duplicate symbol error and is left with no schema fields (DataHub logs a warning and continues). On estates where many topics share common proto types, most protobuf topics after the first will therefore have no schema fields. Enable schema_resolution so those topics fall back to schema inference from message data.

In addition to this, maps are represented as arrays of messages. The following message,

message MessageWithMap {
map<int, string> map_1 = 1;
}

becomes:

message Map1Entry {
int key = 1;
string value = 2/
}
message MessageWithMap {
repeated Map1Entry map_1 = 1;
}

Troubleshooting

If ingestion fails, validate credentials, permissions, connectivity, and scope filters first. Then review ingestion logs for source-specific errors and adjust configuration accordingly.

Schema Parsing Errors

DataHub automatically handles schema parsing errors gracefully and continues processing.

Avro Binary Encoding Errors

Error: avro.errors.InvalidAvroBinaryEncoding: Read 0 bytes, expected 1 bytes

Enable schema resolution to automatically infer schemas:

source:
type: kafka
config:
schema_resolution:
enabled: true
offset_reset_strategy: "hybrid"
Protobuf Duplicate Symbol Errors

Error: Couldn't build proto file into descriptor pool: duplicate symbol

DataHub logs warnings and continues processing. Topics with schema conflicts will use inferred schemas if schema resolution is enabled.

Schema Registry Connection Issues

If schema registry is unavailable, enable schema resolution as a fallback:

source:
type: kafka
config:
connection:
schema_registry_url: "http://localhost:8081"
schema_resolution:
enabled: true

Performance Optimization

For large Kafka clusters with many topics:

source:
type: kafka
config:
topic_patterns:
allow: ["prod_.*", "analytics_.*"]
deny: [".*_temp", ".*_test"]

profiling:
enabled: true
max_workers: 20 # controls both profiling and schema resolution parallelization
sample_size: 100
nested_field_max_depth: 10

schema_resolution:
enabled: true
sample_timeout_seconds: 1.0

Memory Issues with Large Topics

For topics with large messages or high volume, reduce the sample size and recursion depth:

profiling:
enabled: true
sample_size: 50
nested_field_max_depth: 2

Code Coordinates

  • Class Name: datahub.ingestion.source.kafka.kafka.KafkaSource
  • Browse on GitHub
Questions?

If you've got any questions on configuring ingestion for Kafka, feel free to ping us on our Slack.

💡 Contributing to this documentation

This page is auto-generated from the underlying source code. To make changes, please edit the relevant source files in the metadata-ingestion directory.

Tip: For quick typo fixes or documentation updates, you can click the ✏️ Edit icon directly in the GitHub UI to open a Pull Request. For larger changes and PR naming conventions, please refer to our Contributing Guide.