Skip to content
Blume
Esc
↑↓navigate↵open⌘Jpreview
Guides

API reference

Document Kafka events with AsyncAPI

Describe your Kafka topics in AsyncAPI, publish event-driven API documentation with a page per operation, and pair it with producer and consumer code you've run.

By 12 min read

To document Kafka events, write an AsyncAPI 3 document for the service that owns them. Each Kafka topic is a channel whose address is the topic name, each event is a message with a payload schema and a key, and each operation says whether the service sends to that topic or receives from it. Blume's asyncapi() adapter turns the file into event-driven API documentation with a page per operation, and a page you write beside it shows producer and consumer code you've run against a real broker.

The example is a fictional Acme Orders service: it publishes an orders.created event, consumes cancellation requests from orders.cancel-requests, and runs on a production cluster that requires SASL. You need a Blume project to add it to; the OpenAPI guide covers starting one.

Two limits up front. A browser can't open a Kafka connection, so Kafka pages get a message composer and copyable kcat commands, not a Send button. And payloads described in Avro or Protobuf show a note instead of field tables. If all you need is one standalone HTML page of the contract, AsyncAPI's own @asyncapi/html-template generator covers that without a docs site.

Decide what the document describes

An AsyncAPI 3 document describes one application, and each operation's action is written from that application's side: send means it produces to the channel, and receive means it consumes from it. The people reading your reference are on the other side of every operation:

TopicThe Orders serviceAn integratorBadge on the page
orders.createdProduces (send)ConsumesSEND
orders.cancel-requestsConsumes (receive)ProducesRECEIVE

Blume's generated samples follow the reader's side: a send page shows a command that consumes, and a receive page shows one that produces. Coming from AsyncAPI 2, note that the words flipped. A 2.x publish meant other applications publish to yours, so it becomes receive, and subscribe becomes send.

Write the AsyncAPI document

Save this as asyncapi.yaml at the project root, beside blume.config.ts:

asyncapi: 3.1.0
info:
  title: Acme Orders events
  version: 1.0.0
  description: |
    The Kafka topics the Acme Orders service writes to and reads from. Every
    message value is JSON, and every message is keyed by order ID, so all
    events for one order land on the same partition, in order.
  tags:
    - name: Orders
      description: Events the Orders service publishes when an order changes.
    - name: Commands
      description: Requests the Orders service consumes from other services.
defaultContentType: application/json

servers:
  production:
    host: kafka.acme.example:9093
    protocol: kafka-secure
    description: Production cluster. Clients connect over TLS and authenticate with SASL/SCRAM-SHA-512.
    security:
      - $ref: "#/components/securitySchemes/scram"
  local:
    host: localhost:9092
    protocol: kafka
    description: A broker on your own machine, for trying the samples. No authentication.

channels:
  ordersCreated:
    address: orders.created
    title: Orders created
    description: One message per order, written after the order is committed to the database.
    messages:
      OrderCreated:
        $ref: "#/components/messages/OrderCreated"
    bindings:
      kafka:
        partitions: 12
        replicas: 3
        topicConfiguration:
          cleanup.policy: ["delete"]
          retention.ms: 604800000
        bindingVersion: "0.5.0"
  cancelRequests:
    address: orders.cancel-requests
    title: Cancel requests
    description: Requests from other services to cancel an order that hasn't shipped.
    messages:
      CancelOrder:
        $ref: "#/components/messages/CancelOrder"
    bindings:
      kafka:
        partitions: 3
        replicas: 3
        bindingVersion: "0.5.0"

operations:
  orderCreated:
    action: send
    channel:
      $ref: "#/channels/ordersCreated"
    title: Order created
    description: |
      The Orders service publishes this event once per order, after the order
      is saved. Delivery is at least once, so a consumer can see the same
      event twice: deduplicate on `eventId`. Ignore fields you don't
      recognize, because new optional fields can appear without a new topic.
    tags:
      - name: Orders
    messages:
      - $ref: "#/channels/ordersCreated/messages/OrderCreated"
  cancelOrder:
    action: receive
    channel:
      $ref: "#/channels/cancelRequests"
    title: Cancel an order
    description: |
      The Orders service consumes this topic in the `orders-service` consumer
      group. Produce one message per request, keyed by the order ID. An order
      that has already shipped isn't canceled, and the request is dropped.
    tags:
      - name: Commands
    bindings:
      kafka:
        groupId:
          type: string
          enum: ["orders-service"]
        bindingVersion: "0.5.0"
    messages:
      - $ref: "#/channels/cancelRequests/messages/CancelOrder"

components:
  securitySchemes:
    scram:
      type: scramSha512
      description: SASL/SCRAM-SHA-512 over TLS. Each service gets its own username and password from the platform team.
  messages:
    OrderCreated:
      name: OrderCreated
      title: Order created
      summary: A new order was placed and saved.
      contentType: application/json
      headers:
        type: object
        required: [eventType]
        properties:
          eventType:
            type: string
            enum: ["order.created"]
            description: Lets a consumer that reads several topics route a message without parsing its value.
      payload:
        $ref: "#/components/schemas/OrderCreatedPayload"
      examples:
        - name: cardOrder
          summary: A two-item order
          headers:
            eventType: order.created
          payload:
            eventId: evt_01J9ZK3Q8X
            orderId: ord_7Hq2
            customerId: cus_19f
            createdAt: "2026-09-27T10:15:00Z"
            currency: USD
            total: 4800
            items:
              - sku: MUG-BLUE
                quantity: 2
                unitPrice: 1400
              - sku: TEE-M
                quantity: 1
                unitPrice: 2000
      bindings:
        kafka:
          key:
            type: string
            description: The order ID, the same value as `orderId` in the payload.
          bindingVersion: "0.5.0"
    CancelOrder:
      name: CancelOrder
      title: Cancel order
      summary: Ask the Orders service to cancel an order.
      contentType: application/json
      payload:
        $ref: "#/components/schemas/CancelOrderPayload"
      examples:
        - name: customerRequest
          summary: A customer asked support to cancel
          payload:
            orderId: ord_7Hq2
            reason: customer_request
            requestedBy: support-desk
      bindings:
        kafka:
          key:
            type: string
            description: The order ID, the same value as `orderId` in the payload.
          bindingVersion: "0.5.0"
  schemas:
    OrderCreatedPayload:
      type: object
      required: [eventId, orderId, customerId, createdAt, currency, total, items]
      properties:
        eventId:
          type: string
          description: Unique per event. Deduplicate on it.
        orderId:
          type: string
          description: The order's ID. Also the message key.
        customerId:
          type: string
        createdAt:
          type: string
          format: date-time
          description: When the order was saved, in UTC.
        currency:
          type: string
          pattern: "^[A-Z]{3}$"
          description: ISO 4217 currency code.
        total:
          type: integer
          minimum: 0
          description: Order total in the currency's minor unit, like cents.
        items:
          type: array
          minItems: 1
          items:
            $ref: "#/components/schemas/LineItem"
    LineItem:
      type: object
      required: [sku, quantity, unitPrice]
      properties:
        sku:
          type: string
        quantity:
          type: integer
          minimum: 1
        unitPrice:
          type: integer
          minimum: 0
          description: Price of one unit in the currency's minor unit.
    CancelOrderPayload:
      type: object
      required: [orderId, reason, requestedBy]
      properties:
        orderId:
          type: string
        reason:
          type: string
          enum: [customer_request, payment_failed, out_of_stock]
        requestedBy:
          type: string
          description: The service or team asking, for the audit log.

The .example hosts don't exist. Point them at your brokers.

Servers and security

production uses the kafka-secure protocol and lists a scramSha512 security scheme, so every operation page gets an Authorization section showing scram as SASL/SCRAM-SHA-512, with the scheme's description. Blume labels the other Kafka schemes too (plain, scramSha256, gssapi, X509). It documents them only: the reference never asks readers for broker credentials. The local server gives readers a second target in Try it.

Channels are topics

A channel's address is the topic name. It's what the page heading shows and what the samples connect to. Partition count, replicas, and retention go in the channel's kafka binding, which the page lists under Channel bindings. Skip the binding's own topic field unless it matches the address, since the samples always use the address.

Messages carry the contract

Each message has a payload schema, a headers schema, and a kafka binding whose key documents the message key. The first entry under examples becomes the example payload on the page and the starting value in Try it. Without one, Blume samples a value from the schema, which is valid but tells readers less.

Operations say who does what

An operation's title is its page title and sidebar label, its first tag picks the sidebar group, and its description opens the page. Put the rules a client must follow in the description: deduplicate on eventId, key by order ID, which consumer group reads the topic. The description is also what agents get. An operation's Markdown copy, the same text llms-full.txt and the MCP server serve, carries the description and the signature, SEND orders.created, but not the schema tables. The groupId in cancelOrder's binding records the Orders service's consumer group, and the page shows it under Operation bindings.

Validate the document

Blume reads the document leniently: it skips what it can't place and warns about it, but it doesn't check the document against the AsyncAPI schema. The AsyncAPI CLI does:

npx @asyncapi/cli@6.2.0 validate asyncapi.yaml

On the document above, it reports the file as valid with no governance issues. Its linter notes when a document isn't on the latest version, 3.1.0, which is why the example uses it. Blume reads any 3.x document as it is. A 2.x document also works, but Blume converts it to 3.0 at build time with @asyncapi/converter, which you install yourself.

Mount the event reference

Import asyncapi() from blume/reference, list it under reference, and add a header tab for it:

import { defineConfig } from "blume";
import { asyncapi } from "blume/reference";

export default defineConfig({
  title: "Acme Docs",
  navigation: {
    tabs: [
      { label: "Docs", path: "/" },
      { label: "Events", path: "/events" },
    ],
  },
  reference: [asyncapi({ spec: "./asyncapi.yaml" })],
});

The reference mounts at /events unless you pass route. It never adds a tab on its own, so the Events tab is what makes it reachable, and it also scopes the sidebar to the event pages. Start the dev server with npm run dev. It doesn't watch a spec at the project root, so restart it after each edit to asyncapi.yaml.

Each operation gets its own URL:

OperationPageSidebar group
orderCreated/events/orders/order-createdOrders
cancelOrder/events/commands/cancel-orderCommands

The path is the first tag, then the operation's key under operations, split at camelCase boundaries. An untagged operation is grouped under its channel address instead. Renaming a tag or an operation key moves the page, so add a redirect when you do.

Read the rendered reference

/events is an overview: the title and description from info, the version, both servers as kafka-secure://kafka.acme.example:9093 and kafka://localhost:9092, and a section per tag listing its operations. The Order created page has these parts:

  1. Title and description. Order created, the operation's description, then a SEND badge beside orders.created.
  2. Authorization. scram, labeled SASL/SCRAM-SHA-512 and required, with the scheme's description. It says required even though local has no security: the section lists every scheme the channel's servers declare.
  3. Message. The application/json content type, a payload table with each field's type, whether it's required, its description, and constraints like the currency pattern, then a headers table.
  4. Message bindings. The key schema from the message's kafka binding.
  5. Channel bindings. partitions, replicas, and topicConfiguration.
  6. Side panel. A collapsed Try it composer, a kcat tab under Example, and the example payload under Message.

Cancel an order has the same parts under a RECEIVE badge, plus Operation bindings with the orders-service group ID.

The generated samples

The kcat sample comes from the operation's action and the channel's first server. For the two operations above, it's:

# Order created (SEND): watch the topic
kcat -b 'kafka.acme.example:9093' -t 'orders.created' -C

# Cancel an order (RECEIVE): produce a request
echo '{"orderId":"ord_7Hq2","reason":"customer_request","requestedBy":"support-desk"}' | kcat -b 'kafka.acme.example:9093' -t 'orders.cancel-requests' -P

Open Try it to change them. Pick the local server and the broker becomes localhost:9092. Edit the payload and the echo line follows, while the editor checks your JSON against the payload schema as you type. What Try it doesn't do on a Kafka page is connect: the panel says live sending isn't possible from a browser and leaves the command to you. Only WebSocket operations get a live connection.

The samples carry no credentials and no message key. Against production, add the SASL settings that kcat passes through to librdkafka, and key the message with -K as the contract requires:

echo 'ord_7Hq2|{"orderId":"ord_7Hq2","reason":"customer_request","requestedBy":"support-desk"}' \
  | kcat -b kafka.acme.example:9093 -t orders.cancel-requests -P -K '|' \
    -X security.protocol=SASL_SSL -X sasl.mechanisms=SCRAM-SHA-512 \
    -X sasl.username="$KAFKA_USERNAME" -X sasl.password="$KAFKA_PASSWORD"

Add a producer and consumer page

The reference answers "what's on this topic?" Integrators also need "how do I connect, and what must my code do?" A page in the events folder of your content answers that, and it joins the Events tab beside the generated pages. Keep the client code in files next to it and embed them, so the page shows code you run rather than code you retyped.

The clients use Confluent's Node.js client. Add it to the docs project as a dev dependency, since only the examples use it:

npm install --save-dev @confluentinc/kafka-javascript@1.10.1

Put the three files in docs/events/_examples/. The underscore keeps that folder out of routing, search, and the sidebar. The first file holds the connection settings: a local broker by default, and SASL over TLS once credentials are set.

// Local broker by default. Set KAFKA_USERNAME and KAFKA_PASSWORD to use
// SASL/SCRAM-SHA-512 over TLS, as the production cluster requires.
export const config = {
  "bootstrap.servers": process.env.KAFKA_BROKERS ?? "localhost:9092",
  ...(process.env.KAFKA_USERNAME && {
    "security.protocol": "SASL_SSL",
    "sasl.mechanisms": "SCRAM-SHA-512",
    "sasl.username": process.env.KAFKA_USERNAME,
    "sasl.password": process.env.KAFKA_PASSWORD,
  }),
};

The consumer reads orders.created in its own consumer group, so it gets every event no matter who else consumes the topic, and it skips any eventId it has already processed:

import { KafkaJS } from "@confluentinc/kafka-javascript";
import { config } from "./kafka-config.mjs";

const consumer = new KafkaJS.Kafka().consumer({
  ...config,
  "group.id": "fulfillment",
  "auto.offset.reset": "earliest",
});

// Delivery is at least once. Keep processed IDs in your database in
// production; a Set is enough to show the idea.
const processed = new Set();

await consumer.connect();
await consumer.subscribe({ topics: ["orders.created"] });
await consumer.run({
  eachMessage: async ({ message }) => {
    const event = JSON.parse(message.value.toString());
    if (processed.has(event.eventId)) return;
    processed.add(event.eventId);
    console.log(
      `Fulfill ${event.orderId} (key ${message.key}): ${event.items.length} line items`
    );
  },
});

The producer sends one cancel request, keyed by the order ID, so every request for one order lands on the same partition, in order:

import { KafkaJS } from "@confluentinc/kafka-javascript";
import { config } from "./kafka-config.mjs";

const [orderId, reason = "customer_request"] = process.argv.slice(2);
const producer = new KafkaJS.Kafka().producer(config);

await producer.connect();
await producer.send({
  topic: "orders.cancel-requests",
  messages: [
    {
      key: orderId,
      value: JSON.stringify({ orderId, reason, requestedBy: "support-desk" }),
    },
  ],
});
await producer.disconnect();
console.log(`Requested cancellation of ${orderId}`);

The page itself is short. Each <include> embeds a file as a code block, and meta gives it a title:

---
title: Consume and produce order events
description: Connect to the Acme Kafka cluster from Node.js, consume Order created events, and request cancellations.
---

Locally, the clients connect to `localhost:9092` with no credentials. In
production, set `KAFKA_BROKERS` to `kafka.acme.example:9093`, and set
`KAFKA_USERNAME` and `KAFKA_PASSWORD` to your service's credentials.

<include meta='title="kafka-config.mjs"'>./_examples/kafka-config.mjs</include>

## Consume Order created

The Orders service publishes [Order created](/events/orders/order-created).
Consume it in your own consumer group. Delivery is at least once, so
deduplicate on `eventId`.

<include meta='title="consume-orders.mjs"'>./_examples/consume-orders.mjs</include>

## Request a cancellation

The Orders service consumes [Cancel an order](/events/commands/cancel-order).
Send one message per request, keyed by the order ID.

<include meta='title="request-cancel.mjs"'>./_examples/request-cancel.mjs</include>

The page lives at /events/integrate. Its links to operation pages are ordinary links, so blume validate checks them like any other.

Run the clients against a local broker

Before the page tells anyone the code works, run it. Kafka's official image starts a single-node broker on localhost:9092; create both topics in it:

docker run -d --name kafka -p 9092:9092 apache/kafka:4.3.1
docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic orders.created --partitions 12
docker exec kafka /opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic orders.cancel-requests --partitions 3

Start the consumer and leave it running:

node docs/events/_examples/consume-orders.mjs

In a second terminal, stand in for the Orders service: publish the spec's example event twice, keyed by order ID.

event='{"eventId":"evt_01J9ZK3Q8X","orderId":"ord_7Hq2","customerId":"cus_19f","createdAt":"2026-09-27T10:15:00Z","currency":"USD","total":4800,"items":[{"sku":"MUG-BLUE","quantity":2,"unitPrice":1400},{"sku":"TEE-M","quantity":1,"unitPrice":2000}]}'
printf 'ord_7Hq2|%s\nord_7Hq2|%s\n' "$event" "$event" \
  | kcat -b localhost:9092 -t orders.created -P -K '|'

Among the client's own connection logs, the consumer prints one line, not two, because the second copy has an eventId it has already seen:

Fulfill ord_7Hq2 (key ord_7Hq2): 2 line items

Then send a cancel request and read it back:

node docs/events/_examples/request-cancel.mjs ord_7Hq2
kcat -b localhost:9092 -t orders.cancel-requests -C -e -f '%k %s\n'

kcat prints the key, then the payload, in the shape the contract describes. Run the same steps in CI to catch a client upgrade or topic change that breaks the embedded files.

Build and check

Build the site and serve the production build:

npm run build
npx blume preview

Open the URL it prints and check:

  • Operation URLs. /events/orders/order-created loads on its own, not only from the sidebar.
  • Search. Searching "order created" finds the operation. Search matches an operation's title, description, tag, and signature, not the fields in its tables, so name the fields readers look for in the description.
  • The Markdown copy. /events/orders/order-created.md returns the description and SEND orders.created.

Then run npx blume validate to check every internal link, including the ones from integrate.mdx to the operation pages.

Troubleshooting

The build fails with BLUME_ASYNCAPI_UNAVAILABLE

Blume couldn't load the spec, and the message says why: a wrong path, a file with no asyncapi version field, or a 2.x document without @asyncapi/converter installed, in which case the diagnostic suggests the install command. blume dev only warns and leaves the reference out, so a working dev server can hide it.

An operation is missing

Look for BLUME_ASYNCAPI_SKIPPED_OPERATION in the output. It names the operation and the reason: its action isn't send or receive (a leftover 2.x publish, say), or its channel points at a channel that isn't declared under channels.

The reference is empty

BLUME_ASYNCAPI_EMPTY means the file parsed but has no operations. Pages come from operations, not channels, so a document that only declares channels renders nothing. Add an operation per topic your service sends to or receives from.

SEND and RECEIVE look backwards

The badges describe your service, not the reader. If the page for a topic you produce says RECEIVE, the action is flipped in the spec, often from a hand-converted 2.x publish. Change it to send, and the sample switches from a produce command to a consume command too.

The payload shows a note instead of a table

The payload is a multi-format schema, { schemaFormat, schema }, whose schemaFormat is Avro or Protobuf. Blume renders AsyncAPI, JSON Schema, and OpenAPI schemas as tables, and shows a note for other formats. The example payload still appears if the message has examples, so add one, and describe the key fields in the message's description.

The sample uses the wrong broker or topic

The sample uses the first server the channel is available on, and the channel's address as the topic. List the server readers should see by default first, or give the channel its own servers list. If the address differs from bindings.kafka.topic, make them match: the samples never read the binding.

Next step

Validate your event contract

Run the AsyncAPI validator on your document, then add it to the reference and build.

npx @asyncapi/cli@6.2.0 validate asyncapi.yaml
Read the AsyncAPI reference docs

A step here not working for you? Report a broken step.

Keep going.More guides.

Upgrade your docs with Blume.

Install today and ship a production-grade docs site in minutes. Free and open source, forever.

npx blume init