Evolving event-sourced domain events in Go with Avro

Go
Event sourcing
DDD
Avro

In an event-sourced application, the events are the state. A domain event such as SubscriptionStarted is written once, never updated, and read again every time an aggregate is rebuilt or a projection is replayed. The code that reads it keeps changing, so that code has to read every shape the event ever had.

This article shows one way to do that in Go with Avro: derive the Avro schema from the domain event struct, check every change against every stored version before it ships, and let Avro convert old events to the current shape when they are read. The same schemas generate event contracts, Go types that let a projection read other bounded contexts' events without importing their code. Avro handles the shape of an event, not its meaning; the last section shows where that line falls.

The example is a small subscription business with two bounded contexts, Subscriptions and Invoicing, and a customer account overview that combines events from both.

Migrating events is not migrating tables

A relational database and an event store both have to survive schema changes, but they do it in opposite ways.

In a relational database:

  • There is one current shape. A migration (ALTER TABLE, often followed by an UPDATE to backfill) rewrites the existing rows, and the old shape is gone.
  • Backward compatibility is a temporary concern. During a rolling deploy, the previous version of the application must keep working against the new schema, which is why columns are added first and removed later. Once the deploy is done, nothing reads the old shape again.

In an event store:

  • Stored events are not rewritten: an event records what happened, and rewriting it changes history. The "migration" happens at read time instead, on every read.
  • Backward compatibility is permanent. In Avro's terms, it means the current schema can read data written with every older schema. Every version ever written stays in the store, so every change has to be checked against all of them, forever.
  • Forward compatibility, the reverse (old code reading new data), is the rolling-deploy concern here, and it comes back at the end.

The tools map closely:

  • ADD COLUMN currency text NOT NULL DEFAULT 'EUR' becomes a new field with an Avro default.
  • RENAME COLUMN becomes an alias on the field.
  • A backfill UPDATE with business logic becomes an upcaster (a function that converts an old event into the current shape when it is read) or a new event type.
  • Numbered migration files and a runner in CI become a folder of numbered schema versions and a compatibility check in CI.

Some teams do rewrite an event store, copying every event into a new store and transforming it on the way, but that is a heavy operation kept for changes nothing else can handle. Another common choice is JSON events with a version number and hand-written upcasters, which suits events that change rarely. Unless the team writes tests for it, though, nothing checks before deployment that the upcasters still cover every stored version. What draws me to Avro is that its compatibility rules are part of the format, so a machine can check them.

Deriving the Avro schema from the domain event

The first version of the Subscriptions events is plain Go:

package subscriptions

type Address struct {
	Street  string
	City    string
	Country string
}

type Customer struct {
	Name    string
	Address Address
}

type SubscriptionStarted struct {
	ID       string
	Customer Customer
	Plan     string
}

type SubscriptionCancelled struct {
	ID     string
	Reason string
}

var Events = []any{SubscriptionStarted{}, SubscriptionCancelled{}}

The examples use struct-avro to derive schemas and ln80/avro, a fork of hamba/avro, to encode. Every event is stored in an envelope whose Data field becomes an Avro union of the namespace's event types. The helper below builds that envelope schema from sample events; a sample's field values become the Avro defaults:

package envelope

import (
	"reflect"

	avro "github.com/ln80/avro/v2"
	stravro "github.com/ln80/struct-avro"
)

type Envelope struct {
	ID       string `avro:"ID"`
	StreamID string `avro:"StreamID"`
	Type     string `avro:"Type"`
	At       int64  `avro:"At"`
	Data     any    `avro:"Data" ev:",inject=events"`
}

func Schema(api avro.API, namespace string, samples ...any) (*avro.RecordSchema, error) {
	entries := make([]stravro.RecordEntry, 0, len(samples))
	for _, s := range samples {
		t := reflect.TypeOf(s)
		entries = append(entries, stravro.RecordEntry{
			Name:   namespace + "." + t.Name(),
			Type:   t,
			Sample: s,
		})
	}
	return stravro.BuildEnvelopeSchema(api, stravro.EnvelopeConfig{
		Namespace: namespace,
		RootType:  reflect.TypeOf(Envelope{}),
		RootName:  "Envelope",
		InjectKey: "events",
		Entries:   entries,
		Validate:  true,
	})
}

Avro bytes cannot be decoded without the schema they were written with, so each stored event starts with that schema's ID, a 16-character fingerprint. reg below is a schema registry backed by a folder of schema files. It only accepts a schema that is already stored there, which the job in a later section takes care of before deployment:

id, schema, _, err := reg.GetCurrent(ctx, nil)
if err != nil {
	return err
}
b, err := reg.API().Marshal(schema, evt)
if err != nil {
	return err
}
b, err = reg.AppendSchemaID(b, id)

The first SubscriptionStarted, for a customer in Berlin on the pro plan, is stored as 128 bytes starting with a0c8d03b4ffd12a0.

Reading an old event with a new schema

Later, the business renames plans to tiers and starts billing in more than one currency. Version 2 renames Plan to Tier, keeping the old name as an alias, and adds a Currency whose default, set through the sample, is euros:

type SubscriptionStarted struct {
	ID       string
	Customer Customer
	Tier     string `ev:",aliases=Plan"`
	Currency string
}

var Events = []any{
	SubscriptionStarted{Currency: "EUR"},
	SubscriptionCancelled{},
}

In the generated schema, the two fields become:

[
	{"name": "Tier", "aliases": ["Plan"], "type": "string", "default": ""},
	{"name": "Currency", "type": "string", "default": "EUR"}
]

Reading an event involves two schemas: the writer schema, which the event was encoded with, and the reader schema, which the current code expects. The code extracts the schema ID, fetches the writer schema, and decodes. Because the registry was set up with v2 as the current schema, GetSchema returns the writer schema already resolved against it:

id, data, err := reg.ExtractSchemaID(b)
if err != nil {
	return err
}
_, schema, _, err := reg.GetSchema(ctx, id)
if err != nil {
	return err
}
var evt envelope.Envelope
if err := reg.API().Unmarshal(schema, data, &evt); err != nil {
	return err
}
fmt.Printf("%+v\n", *evt.Data.(*subscriptions.SubscriptionStarted))
// {ID:sub-42 Customer:{Name:Acme GmbH Address:{Street:Hauptstr. 1 City:Berlin Country:DE}} Tier:pro Currency:EUR}

The event written with Plan: "pro" comes back with Tier: "pro", and the currency it never had comes back as EUR. No conversion code was written.

That default carries the same assumption as DEFAULT 'EUR' in a table: every subscription started before version 2 was billed in euros. A default is a statement about the past, which makes it a business decision. If the past is mixed, no default is correct.

Checking every change against every stored version

Resolution only works if the new schema can read every old one, so schema changes go through a job that runs before deployment, for example in CI:

registry := fs.NewDirAdapter("./registry")
contracts := fs.NewDirAdapter("./contracts")

err := tool.NewJobExecuter(&tool.DefaultPrinter{Output: os.Stdout, Err: os.Stderr}).
	GenerateSchemas(build).
	CheckCompatibility(registry).
	PersistSchemas(registry, registry).
	EmbedSchemas(registry, contracts, tool.EmbedConfig{Out: "./contracts"}).
	Execute(ctx)

build returns the current envelope schema of each namespace, using the helper above. The job then:

  1. Checks it against every stored version of its namespace, not only the latest, and stops at the first one it cannot read.
  2. Stores it as a new version, unless a schema with the same fingerprint is already stored.
  3. Generates the event contracts described in the next section.

The fingerprint ignores defaults, so the library adds a field derived from the defaults and pii annotations to the envelope. Changing only a default still produces a new version.

After version 3 adds a Region to the address, which gets an empty-string default like every other field, the registry holds three versions of Subscriptions and one of Invoicing:

registry/
  invoicing/
    1@6b0c76af6947e4e2.json
  subscriptions/
    1@a0c8d03b4ffd12a0.json
    2@aaf5654a3d62313c.json
    3@b984874c7d8d804e.json

Version 4 turns Tier into a structure:

type Tier struct {
	Name  string
	Seats int
}

type SubscriptionStarted struct {
	ID       string
	Customer Customer
	Tier     Tier `ev:",aliases=Plan"`
	Currency string
}

The job refuses it, and nothing is stored:

2. CheckCompatibility
Started...
CheckCompatibility: error: namespace: subscriptions, incompatible with version: 1, err: reader union lacking writer schema record

The message names the envelope's union because the check compares whole envelopes. The underlying reason is that a string written by version 1 cannot be read as a record.

Event contracts for projections in other bounded contexts

The customer account overview needs events from both Subscriptions and Invoicing. Importing their domain packages would couple it to their internals, and is not possible when they live in separate applications or modules. Instead, the job's last step generates, from the registry, a package of event contracts: Go types matching the latest schema of each namespace, plus every schema version, embedded in the binary. For Subscriptions:

// Code generated by ln80/struct-avro. DO NOT EDIT

// SubscriptionStarted is a generated struct.
type SubscriptionStarted struct {
	ID       string   `avro:"ID"`
	Customer Customer `avro:"Customer"`
	Tier     string   `avro:"Tier"`
	Currency string   `avro:"Currency"`
}

In DDD terms, the registry acts as a Published Language: a shared format that contexts use to talk to each other.

One detail matters here. Decoding an old event with its writer schema alone runs without error, but fills the contract types by field name, ignoring aliases and defaults. In the example, that showed the Berlin customer with an empty tier and currency. The projection has to resolve each event against the latest schema of its namespace, found by walking the embedded versions:

for _, b := range storedEvents {
	id, data, err := wf.ExtractSchemaID(b)
	if err != nil {
		return err
	}
	def, err := schemas.Get(ctx, id)
	if err != nil {
		return err
	}
	writer := avro.MustParse(def).(*avro.RecordSchema)
	schema, err := compat.Resolve(latest[writer.Namespace()], writer)
	if err != nil {
		return err
	}

	var evt envelope.Envelope
	if err := api.Unmarshal(schema, data, &evt); err != nil {
		return err
	}
	switch e := evt.Data.(type) {
	case *subscriptions.SubscriptionStarted:
		overviews[e.ID] = &Overview{Customer: e.Customer.Name, Tier: e.Tier, Currency: e.Currency}
	case *invoicing.InvoiceIssued:
		overviews[e.SubscriptionID].Invoiced += e.AmountCents
		// InvoicePaid, not shown, adds the invoice amount to Paid.
	}
}

Over four events (the v1 subscription, an invoice issued and paid, and a subscription written with v3), the overview is:

sub-42 {Customer:Acme GmbH Tier:pro Currency:EUR Invoiced:4900 Paid:4900}
sub-43 {Customer:Globex Inc. Tier:team Currency:USD Invoiced:0 Paid:0}

Event contracts mirror the domain events one to one, which exposes each context's internal events to its consumers. I find that acceptable when the consumers are read-only projections in the same organization, released together with the registry. When they are other teams or external systems, integration events designed for them, free to differ from the domain model, protect both sides better.

What schema evolution does not solve

Avro answers one question: can this old event still be decoded? It does not answer whether the result is still true. Each case below was run against the example.

Avro resolution handles these, and the check accepts them:

  • Adding a field, including in a nested record.
  • Removing a field.
  • Renaming a field or an event type, with an alias.
  • Adding an event type.
  • Widening a type: int to long, float to double, string to bytes.

The check refuses these:

  • Narrowing a type, such as long to int.
  • Changing a field to a type Avro cannot convert, such as a string to a record.
  • Removing an event type, even one no stored event uses, because the check cannot know that.

And some changes are out of Avro's reach:

  • A change of meaning, such as AmountCents starting to hold euros. The schema does not change, so the check sees nothing.
  • Splitting or merging events, such as replacing SubscriptionStarted with SubscriptionCreated and PlanSelected.
  • A default that is not true for every past event, such as a currency in a business that already billed in two.
  • Forward compatibility. A projection whose contracts are older than an event's schema fails with schema not found until it is redeployed with refreshed contracts.

One case needs a warning. Because every field gets a default, renaming a field without the alias also passes the check: Avro sees a removed Plan and a new, empty Tier, and old events decode as {ID:sub-42 Tier:}. Nothing fails; the plan is just gone, so a rename still needs a reviewer who asks for the alias.

For changes Avro cannot express, the options are an upcaster or a new event type. An upcaster keeps one event type in the code, at the cost of conversion code to maintain. A new event type keeps the old one, still handled. I usually prefer the new event type in an event store, because the history then keeps saying what the system recorded at the time.

Where this approach fits

In short:

  • Derive the Avro schema from the domain event struct, and check every change against every stored version before deployment.
  • Let Avro handle additive changes and renames, and use a new event type or an upcaster for changes of meaning.
  • Give projections in other contexts generated event contracts, and resolve each event against the latest contract.

This fits best when producers and consumers are Go services in the same organization and a build step can run the check. It fits less well when consumers expect JSON they can read without a schema, when they are other companies, or when event shapes change so often that the number of versions becomes the problem. If you handle event versioning differently, or see where this approach breaks down, I would like to hear about it.