Implementation:Datahub project Datahub AvroSerializer
| Knowledge Sources | |
|---|---|
| Domains | Java_SDK, Metadata_Emission, Kafka |
| Last Updated | 2026-02-10 00:00 GMT |
Overview
A serializer that converts DataHub MetadataChangeProposal objects into Avro GenericRecord instances for Kafka emission.
Description
AvroSerializer is responsible for transforming metadata change proposals into Avro-encoded records suitable for publishing to Kafka topics. On construction, it loads the MetadataChangeProposal.avsc schema from the classpath and extracts the nested genericAspect and changeType sub-schemas.
The class provides two overloaded serialize methods:
- One accepting a
MetadataChangeProposalWrapper(the higher-level wrapper), which internally converts it to a rawMetadataChangeProposalvia theEventFormatterusing Pegasus JSON format. - One accepting a raw
MetadataChangeProposal, which populates the AvroGenericRecordfields:entityUrn,aspect(withcontentTypeand binaryvalue),aspectName,entityType, andchangeType.
The Avro schema is loaded once at construction time and reused across all serialization calls.
Usage
This class is used internally by KafkaEmitter to serialize metadata change proposals before publishing them to the Kafka MCP topic. It is not typically used directly by SDK consumers.
Code Reference
Source Location
metadata-integration/java/datahub-client/src/main/java/datahub/client/kafka/AvroSerializer.java
Signature
public class AvroSerializer {
public AvroSerializer() throws IOException
public GenericRecord serialize(MetadataChangeProposalWrapper mcpw) throws IOException
public GenericRecord serialize(MetadataChangeProposal mcp) throws IOException
// Visible for testing
Schema getRecordSchema()
}
Import
import datahub.client.kafka.AvroSerializer;
I/O Contract
Inputs
| Method | Parameter | Type | Description |
|---|---|---|---|
serialize |
mcpw |
MetadataChangeProposalWrapper |
High-level metadata change proposal wrapper |
serialize |
mcp |
MetadataChangeProposal |
Raw Pegasus metadata change proposal |
Outputs
Both serialize methods return an Avro GenericRecord containing the following fields:
| Field | Type | Description |
|---|---|---|
entityUrn |
String | The entity URN as a string |
aspect |
Record | Nested record with contentType ("application/json") and binary value
|
aspectName |
String | Name of the aspect being changed |
entityType |
String | Type of entity (e.g., "dataset") |
changeType |
Enum | The change type (e.g., UPSERT) |
Usage Examples
// Internal usage within KafkaEmitter
AvroSerializer serializer = new AvroSerializer();
GenericRecord record = serializer.serialize(metadataChangeProposal);
// The resulting record is published to Kafka
ProducerRecord<Object, Object> kafkaRecord =
new ProducerRecord<>("MetadataChangeProposal_v1", entityUrn, record);
producer.send(kafkaRecord);
Related Pages
- Implementation:Datahub_project_Datahub_KafkaEmitter -- Primary consumer of AvroSerializer
- Implementation:Datahub_project_Datahub_KafkaEmitterConfig -- Configuration for the Kafka emitter