Implementation:Datahub project Datahub KafkaEmitter: Difference between revisions
Auto-imported from implementations/Datahub_project_Datahub_KafkaEmitter.md |
Sync from local file |
||
| Line 118: | Line 118: | ||
== Related Pages == | == Related Pages == | ||
* [[Datahub_project_Datahub_KafkaEmitterConfig]] -- Configuration for KafkaEmitter | * [[Implementation:Datahub_project_Datahub_KafkaEmitterConfig]] -- Configuration for KafkaEmitter | ||
* [[Datahub_project_Datahub_AvroSerializer]] -- Serializes MCPs to Avro GenericRecords | * [[Implementation:Datahub_project_Datahub_AvroSerializer]] -- Serializes MCPs to Avro GenericRecords | ||
* [[Datahub_project_Datahub_MetadataWriteCallback_Interface]] -- Callback interface for async notification | * [[Implementation:Datahub_project_Datahub_MetadataWriteCallback_Interface]] -- Callback interface for async notification | ||
* [[Datahub_project_Datahub_S3Emitter]] -- Alternative emitter targeting S3 | * [[Implementation:Datahub_project_Datahub_S3Emitter]] -- Alternative emitter targeting S3 | ||
* [[Datahub_project_Datahub_RestEmitterConfig]] -- Configuration for the REST-based emitter | * [[Implementation:Datahub_project_Datahub_RestEmitterConfig]] -- Configuration for the REST-based emitter | ||
[[Category:Implementations]] | [[Category:Implementations]] | ||
[[Category:Implementations]] | [[Category:Implementations]] | ||
Latest revision as of 10:35, 27 September 2026
| 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
MetadataChangeProposalWrapperto raw MCP viaEventFormatter, then serializes and publishes. - Emit with raw MCP: Serializes via
AvroSerializerand publishes to the configured topic. The entity URN is used as the Kafka message key. - Callback bridging: Wraps the DataHub
Callbackinside a KafkaCallback, mappingRecordMetadataand exceptions intoMetadataWriteResponse. - Connection testing: Uses
AdminClientto 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()returnsFuture<MetadataWriteResponse>wrapping the KafkaRecordMetadataresult.MetadataWriteResponse.successistrueif the Kafka record was acknowledged without exception.testConnection()returnstrueif the Kafka cluster responds within 5 seconds,falseotherwise.
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
- Implementation:Datahub_project_Datahub_KafkaEmitterConfig -- Configuration for KafkaEmitter
- Implementation:Datahub_project_Datahub_AvroSerializer -- Serializes MCPs to Avro GenericRecords
- Implementation:Datahub_project_Datahub_MetadataWriteCallback_Interface -- Callback interface for async notification
- Implementation:Datahub_project_Datahub_S3Emitter -- Alternative emitter targeting S3
- Implementation:Datahub_project_Datahub_RestEmitterConfig -- Configuration for the REST-based emitter