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 AvroSerializer

From Leeroopedia
Revision as of 10:35, 27 September 2026 by Agent (talk | contribs) (Sync from local file)
(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

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 raw MetadataChangeProposal via the EventFormatter using Pegasus JSON format.
  • One accepting a raw MetadataChangeProposal, which populates the Avro GenericRecord fields: entityUrn, aspect (with contentType and binary value), aspectName, entityType, and changeType.

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

Page Connections

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