Transcript
Ian Cooper: What we're going to talk about is the somewhat thorny topic of if you've moved to event-driven architectures, how do you manage all the APIs that you're now generating? That's an about me slide. The top bit says I'm old. The middle bit is just a list of standard stuff, and there's various places that I've worked over time. I like the point at the bottom which is the trouble with this kind of relationship is it suggests there's some expert learner model going on. I don't think I'm any smarter than you. I just had the fortunate opportunity in my career to get involved in speaking on stage and get confident at it. Let's just pretend that we're all capable of doing this. It is. You are all capable of speaking on stage. I'd recommend doing it. I work on an open-source messaging framework for .NET developers called Brighter. Feel free to check us out. We're not going to not be open source. We have a slightly weird arrangement now whereby you retain the IP and just license us in perpetuity. We can't do a rug pull, just because a few of those have happened recently.
Definition of Terminology
Let's talk about the three pillars, the three things that you're really going to need to think about when you're managing asynchronous APIs. First of all, just a couple of definitions to make sure we're on the same page. We'll talk a bit about endpoints. What are endpoints? Endpoints are places where messages are sent or received. Generally, what we want to do is define for an endpoint what are we going to exchange there. A better way of thinking about that is where messages should be sent, how they should be sent, and what the message should look like. Diagrammatically you can think about it this way. If I have a service called restaurant management that has an endpoint, and at that endpoint I am giving you a message that the restaurant hours have changed that you can listen to. Then, effectively, I need to have a message, data I send or receive.
About that message I'm going to need to agree a couple of things. I'm going to need to agree what metadata I have. In other words, what is in the headers of the message? Because unless you know what metadata I'm putting in the message, you won't have much chance of being able to deal with that. For example, I might have a country code telling you what country the store is in. You may be filtering in order to decide which of those messages to deal with. You need to understand what headers I'm going to publish. You're going to need to understand what I've put in the body. What format is it in? Is it just JSON? Is it text? Is it Avro? Is it Protobuf? What do I put in? Because you're not going to be able to process it unless you know that as well. The channel is the name we tend to give to the logical idea of a virtual pipe over which we're sending messages.
Reality may be different in different pieces of middleware. It might be basically a topic in Kafka or a routing key in RabbitMQ. Whatever the particular preference is, a channel is just a way of saying that. The binding we think about as actually if this channel is an abstract idea, what are the specific ideas? In other words, what is the protocol? Where does the server live? In principle, you can have multiple bindings for a channel. In principle, you could expose the same channel on both Kafka and, say, SNS. There's nothing to stop you in theory doing that. It would be very confusing for most people. In theory, that is possible. There's a really straightforward way of remembering this. This diagram, as you can see, has now got very busy. Let's have a straightforward way. That was originally created by Clemens Vasters for WCF, which was a .NET framework back, this was probably about 15 years ago, called the ABCs.
The address, the channel, the logical pipe over which our messages will flow. The binding, protocol-specific details, such as transports and encoding. The contract, so the message that is the headers, and the data, the payload of what we're sending. When you think about endpoints, what you need to describe to other people to use them are the ABCs. What is the address? What is the binding? What is the contract?
If you are a small org, let's say up to about five teams, and you've got an event-driven architecture, and you need to communicate information, it's relatively straightforward. Typically, the surface area is not big enough that you don't know about each other's services. You probably know what messages they publish. You have a good view of what's happening end-to-end. If you don't, you can just at least be able to walk over to somebody's desk and ask them. Now you probably can hop into their Teams room and Slack them, or just get them on whatever form of video conference works for you. Typically, we all understand how the message flow for our particular application works at that scale. We may have whiteboarded it out. We may have got Miro if we're more virtual and Miro'd it out. That's a company that emerged in the pandemic. If we needed infrastructure for that messaging, typically, we can resolve that ourselves.
Either we have some TicketOps process, we raise a Jira ticket. Probably our platform team is small anyway, and they can go and provision the infrastructure we need. Or even more commonly, we'll just rely on the fact that our framework will probably provision the necessary infrastructure for us. A lot of messaging frameworks will let you just say, we'll provision the infrastructure and we'll go away and create the topics that you need. That tends to work ok at that scale. Any choreography that we need where effectively we've got messages flowing between services, relatively easy because the number of participants is small. The trouble is it begins to break down at scale. First of all, you've got no idea if endpoints even exist at scale. Even if you do understand there are endpoints, it may be very unclear who owns them. We may really have no idea what the contracts are for the messages that are being published.
We may not even know if we do publish them who actually consumes them. A typical problem we encounter at scale is that a team who believes these two other teams consume my event, and they're the people I talk to when I make changes. Forget that, for example, there's a whole load of people listening to it in the data side who are using it for analysis. We only really discover that when they change the message schema having agreed with the two teams, not told data, and suddenly data is calling a production incident. It's very typical that we find that people don't know enough about who consumes their messages and people don't know enough about who sends their messages. If we build our infrastructure from some Jira tickets or from self-provisioning from our frameworks, what do we actually do if we have an incident where that goes missing?
Typically, some fat-fingered problem. Most commonly I've seen it be well-intentioned attempts to reduce cost by removing these things that surely nobody is using anymore. How do we actually recover from that? Because the problem may be that if you begin to look at it in enough detail, you'd have to actually just redeploy pretty much everything to figure out somebody must have owned this particular endpoint, and their self-provisioning is what we need to run to actually get everything working. There's no effective way to recover your business in the event effectively that it goes down other than redeploying the world.
The Three Pillars
There are three things we need to solve in order to conquer most of these problems at scale. The first thing we need to do is answer the question, how are we going to discover our endpoints? How are we going to discover our asynchronous event estate, and be able to find where are our endpoints? What messages? What are the contracts? We're going to need governance. We're going to need to make sure that it's possible that we govern how we change and evolve our schemas so that we essentially don't break other parts of our business. We're going to need to think about how we provision. How do we provision all of these topics that we are going to use such that we can effectively recover or understand where we've drifted?
Discovery
Let's take this one by one. First of all, discovery. Before we began to look at this more closely a few years back, the problem we'd have is the most typical way to find a message or an endpoint was to go into GitHub and search across all of the enterprise repos, and try and find an implementation of that message or an implementation of that endpoint, and try and guess what it might be called. That becomes a problem as soon as you need a keyword like restaurant, because you're going to get a lot of hits. That's not particularly working very well. The other option we tried was plaintive messages in Slack. The people posting into Slack, going, "Does anyone know who publishes the restaurant hours changed event? Anyone know which team owns that?" Then someone would be like, "I think Tracy knows who publishes that. We were discussing it the other day." It's this long chain of discovery.
It doesn't sound very, we're on the top of our system here. If you're lucky, there may be a wiki page in something like Confluent. Probably someone updated it about 18 months ago. The last time somebody actually cared enough to go, we should really document all of our messages and endpoints in the team. He's probably moved on. Maybe we promoted him because he was the person that actually did that. Now we're left with a very outdated wiki page that we don't trust anymore. If I'm a producer of message, how am I going to find who consumes stuff? One of the ones that typically comes up for us is someone has ill-advisedly decided to use a decimal for money or an integer for an ID. The number of people that regret making their IDs an integer in the messages. Again, typically, can I ask in Slack, does anyone consume my message?
Can I search in GitHub saying, I wonder if I can find all the users of my message? It turns out really that you've got the problem that nobody in data actually is in the same channels as you. They're not looking for your message for their process that is reading your events to do some analytical report that the CEO looks at every morning.
How do you evolve events in that case? There is prior art here. That prior art is that we solved this problem quite some time ago for our synchronous APIs with things like OpenAPI, GraphQL, even gRPC has its own docs. Why do we not historically embrace the importance of this ABC, address, binding, and contract for asynchronous APIs? The answer came around, and people said, that's what we need to do. That is where AsyncAPI came from. AsyncAPI is part of the Linux Foundation. Fran Mendez was one of the people that started it back in the day. Their goal was literally that. It was say, there's OpenAPI, can't we create a version of that that works the messaging? We're currently, I think, on version 3.1 of AsyncAPI. If you go and look at earlier versions, you can very clearly see the OpenAPI heritage. If you're used to working with OpenAPI, much of the structure of how it works will seem very familiar to you.
It is very accessible for people who've got that kind of workflow. Typically, with AsyncAPI, what do we capture? We capture the server. Server is typically a broker. AsyncAPI also supports WebSockets, so we can't use the term broker. You can think about that as being the broker in most cases. The channel. Remember we were talking about channels earlier, the virtual pipe over which messages flow. It has a notion of a channel. Channels live on a server. Each channel has messages that pass over that channel. Some of our channels will be what we call data type channels. In other words, we restrict them to a single message. That's generally, in my experience, easy to maintain. Some of our channels are going to have multiple messages. That's particularly often common where people are trying to use some strategy of having events that represent change or a delta rather than events that represent a snapshot.
You can have multiple messages against the one channel in the docs. Then the way you do it is for any given application, so anything that sends or receives messages, we declare operations. An operation says we are going to send or receive a message from a channel. Then we can add to that, bindings, where the bindings are protocol-specific information. The binding might be that this is Kafka. There are also information, spoilers for later, like a number of partitions, for example, that that Kafka topic may have.
I'll give you a view of the structure. I'm going to show you an actual example of two or three of our event schemas in AsyncAPI at the end, because it makes it easy just to cut from the slides. I can show you a lot of our stuff at the end. This is a structure. Those of you who are fans of OpenAPI may well be familiar with the idea. We essentially have the Root Object. Then we have an Info Object, which has metadata for this particular application. The most useful things here are things like version. Who is the contact point? What team owns this particular endpoint? Servers Objects, so how do you connect to the server? Very typically, you may have multiple servers here, because you may have QA, staging, and production environments, for example. You can list all of those. The Channels Object. Essentially, what channels are we currently hosting on that server right now?
Channels can be a reference to all these things in the same way OpenAPI does. Rather than having to inline channels on consumers, receivers, and keep repeating it everywhere, you can use the reference syntax, just point it towards a larger document, which effectively has all these reuse definitions along with messages. The channels being used by the object, that will also contain the messages that we intend to send over the channels. An Operation Object, which essentially says, what are we doing? Because we're documenting with one of these, a given application. What do we do to this channel? A Components Object, that is essentially this repository where we can centrally put message schemas, channel definitions, and can live in a separate file. Then, tags. In case we have some tagging schema, a system that we want to use. What we don't show here is there are extensions. They can become quite important to you, you'll discover when we're doing later on. Binding information, which can be added to any of these objects where appropriate. That is the structure of AsyncAPI. Again, I'll show you an actual document at the end. You don't need it right now.
As well as AsyncAPI, obviously, when you represent any one of these messages that you're sending over, AsyncAPI supports a number of formats you can document the message in. It supports JSON schema. That's its default. That is essentially using JSON schema to describe the message. It supports Avro because Kafka, and more importantly, Kafka's ecosystem tends to rely on you having Avro definitions for things. Protobuf, which tends to be the least commonly used in my experience, but places that effectively have a lot of investment in gRPC, already, but tend to quite often use protobuf. What else do we need to worry about? The other thing we might need to worry about is the metadata. AsyncAPI, in the same way it lets you choose your schema format for documenting messages, also doesn't really say anything about what's your header format going to be. You get to choose it.
Typically, people nowadays will tend to use CloudEvents, which is standardized metadata. It has a number of required attributes, so the ID, the source, the spec version, and the type. The source and the type are the most useful things there. Essentially, really pretty much who produces it and what is the name of this event, and which version of the schema are we talking about. CloudEvents also comes with some optional metadata, which can be quite useful. The data content type, in other words, how am I actually representing the body? What is the schema? There's a URI, which tells me I can retrieve schema information. Subject, which qualifies the source. If you have a number of things that the source needs to tell you about, you can qualify it with more information. The time, so what time it was sent. Typically, we tend to use subject for a number of people in this source might have produced this.
I need to give you some more information which tells you which part of the source produced this message. There are various extensions. For those of you, for example, that like the Claim-Check pattern where you say I'm going to offload the body of my content into S3 storage because it's too big. I'll send you the message. You should retrieve it on the far side. The support, for example, for doing Claim-Checks for CloudEvents. The name comes from the fact that when they wrote CloudEvents, they intended it to be used mainly for eventing scenarios, in other words, some kind of Pub/Sub. Not for messaging scenarios, in other words, some kind of command-response. Typically, our experience would be people just use CloudEvents as a standard set of headers for everything.
One of the things to be aware with CloudEvents is that it's useful for metadata so that we have consistent metadata. Routing is quite common via CloudEvents. A number of messaging frameworks will look at the event type that you've sent it and decide where to route that in an application. That's not uncommon. When you have multiple message types on the same channel, what I need to be able to do is to decide how to deserialize that and route it to a handler. That is tricky if the only way I can do that is by inspecting the actual body and trying to figure that out. A typical way we do that is to use CloudEvents to go and look at the type of the event. From the type of the event, we can then determine how to deserialize it and what to route it to to handle the message.
That's a pretty common approach to these multi-event topics. Effectively, it gives us the ability to do things, that self-describing model. CloudEvents supports two ways of dealing with your existing transport called binary and structured. The first thing you should note is that binary is not binary and structured is not very structured. It's a bit like the joke about the Holy Roman Empire. The three things you should know about it is it wasn't holy, it wasn't Roman, and it wasn't an empire. CloudEvents is a bit similar. What binary actually means is the protocol that I am dealing with has native support for headers. Structured means it doesn't, so I'm going to put it all in the body. My CloudEvents headers, I use an envelope format. I have the headers, and then I have the body, and that lets me put headers in there. Actually, today, there are not many protocols that don't have headers.
You're not forced into structure for that reason. Some of them do have a limited number of headers. The canonical example here would be SNS on the AWS platform where you get about 10 message attributes. If you start adding CloudEvents in, you rapidly consume all of your message attributes. If you wanted to use any of them internally on top of that, we use them, for example, for country information, that kind of thing for the message, you rapidly run out of attributes. It's quite typical in that model. It's become very common in the AWS world to use the structured binding for the CloudEvents headers. I can reserve the headers in the attributes for bespoke ones that I need instead. Typically, ones where I want to do inspection of the headers for routing, rather than need them at the destination. If I use structured, I have to get hold of the body.
That's ok at the destination. I just want it for event routing, and that's fine. If I effectively want to do it on, say, an SNS server and create a filter on my SNS server that will filter given events, I need to use message attributes. Actually, that may help me to reserve those for custom and bespoke things that I need.
Let us imagine you go away, and I hope you will, and you say, great, we're going to create AsyncAPI for all of our messaging endpoints. Go you. You say, we're going to have CloudEvents then for how we're going to describe our messages, and we're going to make sure we have very defined schemas. You may say to me, Ian, we've still not solved the discovery problem. This is on the left-hand side. On the left-hand side, this is basically a screenshot from me in VS Code — I'll show you the thing later — looking at our repo that tends to hold most of our message schemas, our AsyncAPI, our message schemas. I'm still back in the same problem. I've now just got hundreds of files here. We have, I think, 600 to 700 of these things. I've got an agreed format, but I've not really necessarily improved my life very much.
Because now I'm just doing a text search locally. Maybe that's good for agents. They can do local text searches. Has it really helped me do discovery? Typically, on top of this, you'll want some additional tooling. One of the poster children for this is EventCatalog written by a lovely guy, David Boyne. EventCatalog, you can load your schemas and your AsyncAPI information into. It effectively gives you a tool around that catalog that lets you visualize the flow of messages in your system, statically rather than dynamically. Very useful. This picture's not so great, but I'll actually show you this one later. Marmot is an open-source catalog. That's a really useful one if you're budget constrained. We use that one, and I'll give you a little demo at the end. Marmot does a similar job. You just load up your AsyncAPI, your schemas. Marmot is also pitched as a data catalog tool, so it will do lineage, for you as well, which can be very useful to you.
Clemens Vasters here of CloudEvents fame, has also been working recently on a thing called xRegistry. xRegistry says, this is all great, but where's the standard for how we register all of these specifications and schemas in a way that is non-proprietary. It's really another section of the talk, but it includes the idea of, I need to register stuff with Confluent Schema Registry. That's, again, a proprietary format. xRegistry is part of the CNCF, and he's been working on. It's intended to be an open format for registering all this information. Left-hand side, you can see there's this core spec where you get resources, which are a bit of a nebulous concept in that version, which are put in groups in a registry. You can see there are extensions, when you look over on the right-hand side, for endpoints. We're just talking about them. Message definitions, schemas. Nearly all of these support AsyncAPI CloudEvents directly.
Typically, though, they have some REST API to actually upload your content. They can either take in content in that form and spit it back out again. They tend to also have their own proprietary format, in the case of xRegistry, its own, Marmot. It's non-proprietary, but it's tied to that open-source product rather than them using the external schema. All of them work with those standards. Really, xRegistry is very on the minute, lo-fi. You can get a server off the xRegistry project, you can just run an API. There's not really any beautiful UIs and tooling around it. You have to build that yourself. Though, nowadays, with agents, how hard is that? Marmot comes with visualization. The EventCatalog is probably the poster child of visualization, if you have the budget.
AsyncAPI, effectively, that's our OpenAPI equivalent for documenting our endpoints. It's where we can say, what is the address? In other words, what are the servers? What are the contracts? What messages flow over the channels on those servers? What are the bindings? What are the protocols that we need to support? CloudEvents, which says, our messages, as well as a schema in the body, they may have metadata. Maybe we should have some broad agreement about how metadata looks so we can be interoperable. A registry of some sort where having got all this information together, I want to move beyond searching through a repo of them and have some way of just doing discovery against all these APIs and schemas that I've created. You'll want a registry. Together, you can answer the question, what events exist? Who produces them? Who consumes them? Or, what do they look like?
Governance
Next of our pillars was governance. How do we ensure schemas are consistent, compatible, and evolving safely? Typically, when we think about services talking to each other, we tend to prefer to have some open host service. It's the Domain-Driven Design term for, we have an API that provides an external contract which differs from potentially the representation we have inside. That is the point of coupling. It's necessary coupling because otherwise we can't communicate. It's typically data coupling. It might be stamp if we basically depend upon fields that we don't use. That's the best form of coupling we can really get away with. That tends to be all right. It's a contract. I say I want to be able to vary my bounded context inside. My implementation details may change because I may think of better ways of solving this problem. Or I may have to acquire new responsibilities which don't necessarily impact you but may result in me changing stuff.
In order to protect you from that internal change, I want to create a contract which I will honor with you for whatever period of time we agree is reasonable where I won't change that. Your service can depend on me safe in the knowledge that I have a contract with you which I will keep honoring and I will definitely tell you if I need to change that contract. It's a necessity in a world where we may have services owned by different teams that effectively we have some notion of a contract between them to ensure there's a reliable service offering. It's this whole old idea of good fences make good neighbors. That contract is something we're going to agree. That makes the coupling which must exist habitable for both of us. We have a number of ways we can think about that coupling. We can be backwards compatible.
That means that a new version of the consumer can read old data. If we say, the consumer's expecting to receive some additional fields. My typical example here is we've suddenly decided to give location data and latitude and longitude. I launch my consumer that can effectively have those fields. I can just default them out to null island until the producer starts producing those new messages. That means we upgrade the consumers first. When the consumers are all upgraded, we start sending the new message. Forwards compatible. Forwards compatible means that I'm going to send you a message. The producer needs to add some data to it. It's all right, the consumers have been written away. That means that they will ignore additional fields that they don't understand. It's forwards compatibility. There's full compatibility, which is we are both backwards and forwards. I can add messages that you don't understand. You can ignore them effectively, and you can default things that might be missing. None, in other words, I don't like any of you. I don't want friends anymore.
Typically, we have a set of rules about how we evolve schemas. Optional fields are safe. Removing fields, you have to deprecate. When all the consumers are deprecated, you can then remove. Renaming is effectively a remove then add. Change types is almost always a breaking change. We've done change types by effectively treating it almost like rename, which is, add a field with a new type. Change everyone over to consuming a new field with a new type. Essentially remove the old type. Occasionally then do a rename operation on the one you wanted to change the type on. Very painful. How are we going to make sure that everybody essentially is going to keep their schemas, their contracts correct. It's a genuine problem. When I first arrived at Just Eat Takeaway, it was the first problem they dumped on my desk was people keep breaking each other. We try really hard at this, but we still have production incidents caused by somebody has changed the schema and they have not told the other people.
Quite often what happens is the problem is to do with people just not thinking about all of the consumers that might possibly consume from them. They've had a limited conversation. They assume everyone is in the loop and they find out that they're not. Schema is often just a piece of documentation. Even if I've loaded it into EventCatalog, that does not mean that I'm protecting anybody. It just means I'm telling them what it is. I could change it. That doesn't help anybody downstream that they can now see the new documentation for the schema when I've just broken them. In the HTTP world, typically people might solve this problem with something like PACT and consumer-driven contracts, but it's really unsatisfactory in the event world. The reason it's a bit unsatisfactory in the event world is those tools don't work very well because there's no single protocol. There's no HTTP which we can then say, I'll stand up an HTTP server and I will use that to do testing. You'd effectively have to fake out every messaging protocol that there is.
More typically, we use a schema registry. A schema registry says, send me all your schemas. Then, to prevent breakage, what you do is you say, when a producer sends a message, I am going to not let them send it if it doesn't match the schema in the registry. The schema in the registry, when it's updated, will enforce compatibility rules. If I say this schema must be forwards or backwards or full compatible, the schema registry will say, you can't make that change if it would be a breaking change for everybody else. The schema says, I will enforce the contract and I will enforce the ways the contract is allowed to vary. Basically, no breaking change. The producer can only produce a message which matches the schema. Therefore, in order to produce a new message, you must have uploaded the schema. We must have agreed that it doesn't break the consumers by effectively, it passes the rules that we've set.
Now you can produce to it. You can also, and we do some of this, you can say, I will do this as part of my CI/CD pipeline. I will effectively say, I will attempt to register your new schema and check that it passes the rules. Then we can run in a non-production environment. We can then check whether or not we fail to break the producer. Typically, we don't want it to fail in production. We'd like it to fail in staging or earlier on, so that you get a decent warning that you have attempted to load something that's going to break. This may come as a surprise to some teams when it happens because they haven't thought about the change that they want to make and how it might actually impact consumers. We are finding this more and more as we begin to leverage things like Kafka for real-time data, particularly being consumed by agents to provide interactivity with customers, that we rely more and more on those pipelines not breaking because somebody introduced a schema change that broke something five steps down the chain.
We've often got Iceberg tables that depend upon reading schema formats. More and more, it has become really important to enforce rules at the schema. You can automate this. If you are creating AsyncAPI, your AsyncAPI is going to have a schema for the message. What you can do is you can say, I could upload that to my schema registry, validate its compatibility, and essentially then potentially only decide to deploy it if effectively it turns out that I'm not breaking the schema. Typically, we may have some schemas in something like component schemas or external refs, and then we just read them out of there. I'll show you a more comprehensive discussion on what JIT does in a little bit.
Postel's Law is what builds the internet, be conservative in what you send, and liberal in what you accept. While the producers should always check, many of them have the option for the consumer to check, and say, when I read the message, check that the message I'm reading conforms to the schema. I am cautious about doing that. If my consumer can interpret that message, it probably should. Even if the message, for some reason, does not conform to the schema. It can help you out in some PIs that effectively the consumer will honor that message, even if technically it's incorrect, because you can push something through if you need to. However, as soon as you start getting into the world of consumers that do things like create Iceberg tables from Flink, you're going to need to enforce the consumer as well. Standardized schema format. Schema review is part of the process. Schema registry integrated into our CI process, and runtime enforcement for producers.
Provisioning
Provisioning. As we said, a lot of people tend to manually provision using tickets, or they get their framework to do it. We typically see drift. No one's quite sure what production looks like, it isn't actually any more according to what we may have in our documents. There's a lot of complexity sometimes about how we name things in production, like topics, particularly if we have multiple environments, particularly if we have whole loads of teams all dealing with things to do with restaurants. How do we disambiguate stuff? It's generally not easy to figure out what is the source of truth. How do we say what production looks like? Disaster recovery, as said earlier, is basically deploy and hope. What we've abandoned is this model, which says, let's raise some Jira ticket, or also the model would say, let's just get the framework to deploy it for us.
We actually do it from the specification. We take that AsyncAPI document. We take the binding information. We also add, because it's extensible, like OpenAPIs, some additional metadata of our own. Essentially, we say, take all that information. We — as I show you, it's a tiny bit complicated — pull it all together and generate what our messaging infrastructure looks like, and then deploy that. I'm talking more generally here. I will talk more specifically in a second.
AsyncAPI has a generator model. You can take that generator, it reads an AsyncAPI spec, and you can write React code that will say, I want to parse this document, and spit something else out. Originally it was built a lot for saying, spitting out HTML versions of documentation. A lot of folks have begun to use it to say, what I want to do is generate output that I'm going to use to do my provisioning, to do my schema registration and the like. You can use that to do infrastructure provisioning. What you do is you read the AsyncAPI spec. You need to extract the binding and the object. For example, I've got a channel. That channel is going to have some information, for example, saying this is Kafka, and I want 50 partitions of this topic. I want to be replicated across three servers. Then we can provision that information using your weapon of choice.
We use Pulumi. We think it's slightly easier for this particular problem. You can use Terraform, whatever it is you prefer. We can also then have this running in the background. When there is drift, it simply recreates the missing parts of the stack. There are specific bindings for a number of things. For my synths I started it, but mostly the team that I work with does it now internally. We contribute to a few of them to produce this binding information because it needs to be more detailed. Not all of the things that AsyncAPI supports are as detailed because we don't use them. You've got to do document generation right. Publish the catalog. Effectively, also what we can do is publish to Marmot or the like, or publish to EventCatalog as part of that whole process. Quick summary of provisioning. Spec-driven, code generation, auto-registration, DR is tractable, and visibility architecture.
AsyncAPI at Just East Takeaway
What do we do? We have about 603 specs, about 2,000 messaging schemas. We have about three AWS regions. We actually have two platforms. We have about 22,000 commits a year and into a repo that manages most of this information. Our pipeline is we take the AsyncAPI YAML. It gets fed not to the generator, actually to a C# application called Mextrapolater. When we started this, generally, devs were being seconded to do this work. Now, as you'll see later, we tend to write stuff more in Go for DevOps. We tend to use C# still for application code. That generates JSON artifacts. What we do is we take all of these specs and we figure out, this channel is used across all of these specs. That's one channel. In other words, one topic, "Here are all the subscribers that I can see for that." Then it figures out stuff like, here are the IAM rules I'm going to need in AWS.
Here are the ACLs I'm going to need in Kafka. We do a whole lot of processing to figure out what does that actually mean or look like in terms of a change that we want to make to production. We then take that, which is a whole lot of JSON, and we zip it up and put it on S3. Then other people can then get it from there and use it to do stuff. We've basically got a Go CLI called mex deploy, mex stands for message exchange, which will take that artifact when it appears and will then work through AWS Confluent in order to deploy it. Register a schema in the schema registry, and it will pull credentials that it needs out of Vault for any secrets that are necessary to run that process. Up on the top, you can also see we talk to Marmot.
We basically say, take this artifact. To Marmot, we're going to give you our information so that you can give us a data catalog that people can browse. We do schema validation in the pipeline. When we get the schema out via the Templater, we say, ok, what is your schema? Let us try and register that with, we use Confluent Schema Registry, given the compatibility that's been agreed. The only real writer to that schema registry is mex. Mex essentially has a token that says, I'm going to change it. It goes and changes it. Then it hands back the token, and devs can't get the token. The only way they can change schemas in the schema registry is by giving us an AsyncAPI application that contains those schemas. We know we have all the bits and pieces that we actually need in order to make this work. Then you can go and look in our data catalog in Marmot and see what is going on.
This is actually a bit of a Marmot catalog. That's showing you tracking created. You can see it's telling us all about that particular event. Then you can go up here, see a whole load of stuff about the various assets that we have. That's Marmot, and it's pretty easy to find stuff and browse inside it. The other thing I can show you just briefly here is, this is a typical AsyncAPI spec with a whole lot of metadata. The reason I want to show you this is because once you get metadata and stuff in, they're a bit more complicated examples, but you can see there's metadata, the channels. You can see the channel has a name. What happens at runtime is we will add stuff to that name to make it unique and disambiguate it. You can see we've got a message. You can see the message declares a payload.
There are references because it goes to a components file, because it's quite commonly shared. There's binding information, which means here we're looking at specific protocols. This is custom binding for us. What environments are we deploying it to? Then it's basically saying this is actually an SNS message. You'll see somewhere else down here, we've probably got Kafka. There's a Kafka binding here that you can see. That's the picture.
The Virtuous Cycle
What's still hard? There's no notion of being able to support consumer-driven contracts. There's just no PACT equivalent. There's no way of seeing when you document this stuff, I've got endpoints to my application that come in, endpoints that go out. What's the relationship between the two? We have no way of defining that so you can't see flow. There are some works by folks who make OpenAPI, and I think it's on a flow-based representation for OpenAPI and they may extend it to AsyncAPI. There's no semantic compatibility here. Just because I have schema doesn't mean that the intent behind those things is actually tied together. What's hard? Adoption. Teams tend to play game theory. Everybody else should document their stuff, I don't need to. Quite often we have to simply say you cannot provision anything apart from via this process because that's the only way to force you to do it. This doesn't do anything with observability, but there were talks about how to do OTel and messaging. Things like xRegistry exist, but it's not clear where they're going as opposed to the more proprietary things like EventCatalog or open-source alternatives like Marmot. They're not all pulling in the same direction.
Takeaways
Treat your async APIs the same way you do your sync APIs, and document them. Specs are the source of truth, and you can use them for code gen, schema registration, infra provisioning. Governance as a pipeline gate and not a wiki page. Schema registries with compatibility modes make evolution safe and auditable. What it gives you is visibility into your event-driven estate, leverage to enforce standards, and control over your infrastructure.
See more presentations with transcripts