Jump to content

Connect SuperML | Leeroopedia MCP: Equip your AI agents with best practices, code verification, and debugging knowledge. Powered by Leeroo — building Organizational Superintelligence. Contact us at founders@leeroo.com.

Implementation:Datahub project Datahub KafkaEmitter

From Leeroopedia
Revision as of 14:43, 16 February 2026 by Admin (talk | contribs) (Auto-imported from implementations/Datahub_project_Datahub_KafkaEmitter.md)
(diff) ← Older revision | Latest revision (diff) | Newer revision → (diff)


Knowledge Sources
Domains Java_SDK, Metadata_Emission, Kafka
Last Updated 2026-02-10 00:00 GMT

Overview

An Emitter implementation that publishes metadata change proposals directly to a Kafka topic, bypassing the DataHub GMS REST API.

Description

KafkaEmitter implements the Emitter interface to emit metadata changes by serializing MetadataChangeProposal objects to Avro and publishing them to a configurable Kafka topic (default: MetadataChangeProposal_v1). This approach writes directly to the Kafka event stream that DataHub GMS consumes, which can be useful for high-throughput ingestion scenarios or when the REST API is not accessible.

The emitter configures a KafkaProducer using settings from KafkaEmitterConfig, including bootstrap servers, Schema Registry URL, and any additional producer or Schema Registry properties. Keys are serialized as strings (entity URNs), and values are serialized using KafkaAvroSerializer from the Confluent library.

Key behaviors:

  • Emit with wrapper: Converts MetadataChangeProposalWrapper to raw MCP via EventFormatter, then serializes and publishes.
  • Emit with raw MCP: Serializes via AvroSerializer and publishes to the configured topic. The entity URN is used as the Kafka message key.
  • Callback bridging: Wraps the DataHub Callback inside a Kafka Callback, mapping RecordMetadata and exceptions into MetadataWriteResponse.
  • Connection testing: Uses AdminClient to list topics with a 5-second timeout.
  • UpsertAspectRequest: Not supported over Kafka; throws UnsupportedOperationException.

Usage

Use KafkaEmitter when you want to bypass the REST API and write metadata changes directly to Kafka. This is appropriate for high-throughput batch ingestion or environments where the GMS REST endpoint is unavailable but Kafka is directly accessible.

Code Reference

Source Location

metadata-integration/java/datahub-client/src/main/java/datahub/client/kafka/KafkaEmitter.java

Signature

public class KafkaEmitter implements Emitter {

    public static final String DEFAULT_MCP_KAFKA_TOPIC = "MetadataChangeProposal_v1";

    public KafkaEmitter(KafkaEmitterConfig config) throws IOException
    public KafkaEmitter(KafkaEmitterConfig config, String mcpKafkaTopic) throws IOException

    public Future<MetadataWriteResponse> emit(
        MetadataChangeProposalWrapper mcpw, Callback datahubCallback) throws IOException
    public Future<MetadataWriteResponse> emit(
        MetadataChangeProposal mcp, Callback datahubCallback) throws IOException
    public Future<MetadataWriteResponse> emit(
        List<UpsertAspectRequest> request, Callback callback) throws IOException

    public boolean testConnection()
        throws IOException, ExecutionException, InterruptedException
    public void close() throws IOException
    public Properties getKafkaConfigProperties()
}

Import

import datahub.client.kafka.KafkaEmitter;
import datahub.client.kafka.KafkaEmitterConfig;

I/O Contract

Inputs

Method Parameter Type Description
Constructor config KafkaEmitterConfig Kafka and Schema Registry connection settings
Constructor mcpKafkaTopic String Optional custom topic name (default: MetadataChangeProposal_v1)
emit mcpw MetadataChangeProposalWrapper High-level metadata change proposal
emit mcp MetadataChangeProposal Raw metadata change proposal
emit datahubCallback Callback Asynchronous completion callback

Outputs

  • emit() returns Future<MetadataWriteResponse> wrapping the Kafka RecordMetadata result.
  • MetadataWriteResponse.success is true if the Kafka record was acknowledged without exception.
  • testConnection() returns true if the Kafka cluster responds within 5 seconds, false otherwise.

Usage Examples

KafkaEmitterConfig config = KafkaEmitterConfig.builder()
    .bootstrap("kafka-broker:9092")
    .schemaRegistryUrl("http://schema-registry:8081")
    .build();

try (KafkaEmitter emitter = new KafkaEmitter(config)) {
    // Test connectivity
    boolean connected = emitter.testConnection();

    // Emit a metadata change
    Future<MetadataWriteResponse> future = emitter.emit(mcpWrapper, new Callback() {
        @Override
        public void onCompletion(MetadataWriteResponse response) {
            System.out.println("Success: " + response.isSuccess());
        }

        @Override
        public void onFailure(Throwable exception) {
            exception.printStackTrace();
        }
    });

    MetadataWriteResponse response = future.get();
}

Related Pages

Page Connections

Double-click a node to navigate. Hold to expand connections.
Principle
Implementation
Heuristic
Environment