Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,9 @@ rows:
- use_case: "[Example: Convert records for a specific topic to Avro](/event-gateway/policies/record-transcode-produce/examples/convert-topic-to-avro-on-condition/)"
description: |
Use a `condition` to convert records produced to a single topic to Avro, and leave records on other topics unchanged.
- use_case: "[Tutorial: Produce Kafka records as Avro](/event-gateway/produce-kafka-records-as-avro-with-event-gateway/)"
description: |
Let a producer send plain JSON, and store records as Avro in the backend cluster using a Confluent Schema Registry.
{% endtable %}
<!--vale on-->

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ prereqs:
- title: Install kafkactl
position: before
include_content: knep/kafkactl
- title: Start a local Kafka cluster
- title: Start a local Kafka cluster and schema registry
position: before
include_content: knep/docker-compose-start

Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,350 @@
---
title: Produce Kafka records as Avro with {{site.event_gateway}}
content_type: how_to
breadcrumbs:
- /event-gateway/

permalink: /event-gateway/produce-kafka-records-as-avro-with-event-gateway/

products:
- event-gateway

works_on:
- konnect

tags:
- event-gateway
- kafka
- confluent
- schema-registry
- avro

description: "Let producers send plain JSON, and store their records as Avro in the backend cluster using a Confluent Schema Registry."

tldr:
q: How can producers send JSON while records are stored as Avro?
a: |
1. Create a Schema Validation policy (produce phase) to parse a producer's JSON records.
1. Nest a Record Transcode Produce policy that converts the parsed record to Avro using a Confluent Schema Registry.

tools:
- konnect-api

prereqs:
inline:
- title: Install kafkactl
position: before
include_content: knep/kafkactl
- title: Start a local Kafka cluster and schema registry
position: before
include_content: knep/docker-compose-start

cleanup:
inline:
- title: Clean up {{site.event_gateway}} resources
include_content: cleanup/products/event-gateway
icon_url: /assets/icons/gateway.svg

min_version:
event-gateway: '1.3'

related_resources:
- text: Record Transcode Produce policy
url: /event-gateway/policies/record-transcode-produce/
- text: Schema Validation policy
url: /event-gateway/policies/schema-validation-produce/
- text: Schema Registry entity
url: /event-gateway/entities/schema-registry/
- text: Consume Kafka records as JSON with {{site.event_gateway}}
url: /event-gateway/consume-kafka-records-as-json-with-event-gateway/
- text: Validate Avro messages with Confluent Schema Registry
url: /event-gateway/validate-avro-messages-with-schema-registry/
---

## Overview

In this guide, you'll learn how to let producers send plain JSON records while {{site.event_gateway_short}} stores them as Avro in the backend cluster.

Not every producer wants to run a Schema Registry client or serialize Avro directly.
The {{site.event_gateway_short}} [Record Transcode Produce policy](/event-gateway/policies/record-transcode-produce/) converts an already schema-validated record to a different format on its way to the backend cluster, so this producer can keep sending JSON while the cluster and every other consumer of the topic only ever see Avro.

We'll use an `orders` topic that holds order records. A producer that connects through the virtual cluster sends plain JSON, and {{site.event_gateway_short}} converts each record to Avro before it reaches the cluster.

Here's how the data flows through the system:

{% mermaid %}
flowchart LR
P[Producer<br/>JSON] --> SV

subgraph produce [Event Gateway Produce policy chain]
SV[Schema <br>Validation<br/>Parse JSON] --> RT[Record Transcode<br/>convert to Avro]
end

RT --> K[Kafka <br>Broker<br/>Avro records]
{% endmermaid %}

## Create a backend cluster

{% include knep/create-backend-cluster.md insecure=true %}

## Create a virtual cluster

Create a virtual cluster that the producer connects to:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/virtual-clusters
status_code: 201
method: POST
body:
name: orders_vc
destination:
id: $BACKEND_CLUSTER_ID
dns_label: orders
authentication:
- type: anonymous
acl_mode: passthrough
extract_body:
- name: id
variable: VIRTUAL_CLUSTER_ID
capture:
- variable: VIRTUAL_CLUSTER_ID
jq: ".id"
{% endkonnect_api_request %}
<!--vale on-->

## Create a listener with a forwarding policy

Create a [listener](/event-gateway/entities/listener/) to accept connections:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/listeners
status_code: 201
method: POST
body:
name: orders_listener
addresses:
- 0.0.0.0
ports:
- 19092-19105
extract_body:
- name: id
variable: LISTENER_ID
capture:
- variable: LISTENER_ID
jq: ".id"
{% endkonnect_api_request %}
<!--vale on-->

Create a [Forward to Virtual Cluster policy](/event-gateway/policies/forward-to-virtual-cluster/) to forward traffic to the virtual cluster:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/listeners/$LISTENER_ID/policies
status_code: 201
method: POST
body:
type: forward_to_virtual_cluster
name: forward_to_orders_vc
config:
type: port_mapping
advertised_host: localhost
destination:
id: $VIRTUAL_CLUSTER_ID
{% endkonnect_api_request %}
<!--vale on-->

For demo purposes, we're using port mapping, which assigns each Kafka broker to a dedicated port on the {{site.event_gateway_short}}.
In production, we recommend using [SNI routing](/event-gateway/architecture/#hostname-mapping) instead.

## Create a Schema Registry entity

Create a [Schema Registry](/event-gateway/entities/schema-registry/) entity that points to the Confluent Schema Registry running locally.
Since the {{site.event_gateway_short}} data plane runs in the same Docker network as the Schema Registry, use the container hostname `schema-registry`:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/schema-registries
status_code: 201
method: POST
body:
name: local-schema-registry
type: confluent
config:
schema_type: avro
endpoint: http://schema-registry:8081
timeout_seconds: 10
extract_body:
- name: id
variable: SCHEMA_REGISTRY_ID
capture:
- variable: SCHEMA_REGISTRY_ID
jq: ".id"
{% endkonnect_api_request %}
<!--vale on-->

## Register the target Avro schema

Register the Avro schema that {{site.event_gateway_short}} converts records into, under the `orders-value` subject in the Confluent Schema Registry.
The schema defines three fields: `order_id`, `item`, and `amount`:

<!--vale off-->
{% validation custom-command %}
command: |
curl -sS --fail -X POST http://localhost:8081/subjects/orders-value/versions \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d '{"schema": "{\"type\": \"record\", \"name\": \"Order\", \"fields\": [{\"name\": \"order_id\", \"type\": \"string\"}, {\"name\": \"item\", \"type\": \"string\"}, {\"name\": \"amount\", \"type\": \"double\"}]}"}'
expected:
return_code: 0
render_output: false
{% endvalidation %}
<!--vale on-->

The subject name `orders-value` follows Confluent's default [TopicNameStrategy](https://docs.confluent.io/platform/current/schema-registry/fundamentals/serdes-develop/index.html#subject-name-strategy), which uses the pattern `<topic>-value`.
The Record Transcode Produce policy looks up this subject at runtime, so the schema must already exist before a producer sends records.

## Create a Schema Validation policy

Create a [Schema Validation policy](/event-gateway/policies/schema-validation-produce/) that parses a producer's records as JSON during the produce phase.
The Record Transcode Produce policy needs a parsed record, so it must be nested under this policy:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/virtual-clusters/$VIRTUAL_CLUSTER_ID/produce-policies
status_code: 201
method: POST
body:
type: schema_validation
name: validate_json
config:
type: json
value_validation_action: reject
extract_body:
- name: id
variable: SCHEMA_VALIDATION_POLICY_ID
capture:
- variable: SCHEMA_VALIDATION_POLICY_ID
jq: ".id"
{% endkonnect_api_request %}
<!--vale on-->

The `value_validation_action: reject` setting means a batch that holds a record that isn't valid JSON is rejected outright.

## Create a Record Transcode Produce policy

Create the [Record Transcode Produce policy](/event-gateway/policies/record-transcode-produce/) nested under the Schema Validation policy.
Set `output_format: avro` to convert each parsed JSON record to Avro before {{site.event_gateway_short}} writes it to the backend cluster:

<!--vale off-->
{% konnect_api_request %}
url: /v1/event-gateways/$EVENT_GATEWAY_ID/virtual-clusters/$VIRTUAL_CLUSTER_ID/produce-policies
status_code: 201
method: POST
body:
type: transcode
name: convert_orders_to_avro
parent_policy_id: $SCHEMA_VALIDATION_POLICY_ID
config:
failure_mode: reject
output_format: avro
schema_source:
type: reference
schema_registry:
name: local-schema-registry
subject: context.topic.name + '-value'
version: 'latest'
schema_ref_destination:
type: confluent_format
{% endkonnect_api_request %}
<!--vale on-->

Converting to Avro requires a `schema_source`, because Avro needs a schema to serialize the record value.
`subject` and `version` are expressions, so `context.topic.name + '-value'` resolves to `orders-value` for records produced to the `orders` topic, and `'latest'` always resolves to the newest registered version.
`schema_ref_destination: confluent_format` prefixes the converted record with a reference to the schema, so any Confluent-compatible consumer of the topic can deserialize it.

Both policies use `failure_mode: reject`, so a batch that holds a record either policy can't process is never written to the backend cluster.
The alternatives are `passthrough` and `mark`.

## Configure kafkactl

Create a kafkactl configuration with a `vc` context that produces through the virtual cluster with no Schema Registry configured, and a `direct` context that connects straight to Kafka using the Schema Registry for Avro deserialization:

<!--vale off-->
{% validation custom-command %}
command: |
cat <<EOF > kafkactl.yaml
contexts:
vc:
brokers:
- localhost:19092
direct:
brokers:
- localhost:9094
- localhost:9095
- localhost:9096
schemaRegistry:
url: http://localhost:8081
EOF
expected:
return_code: 0
render_output: false
{% endvalidation %}
<!--vale on-->

The `vc` context has no Schema Registry configured, because this producer sends plain JSON through the virtual cluster and doesn't need one.

## Create a topic and produce records

Create the `orders` topic:

<!--vale off-->
{% validation custom-command %}
command: |
kafkactl -C kafkactl.yaml --context direct create topic orders
expected:
message: "topic created: orders"
return_code: 0
render_output: false
{% endvalidation %}
<!--vale on-->

Produce two order records as plain JSON through the virtual cluster:

<!--vale off-->
{% validation custom-command %}
command: |
echo '{"order_id":"1001","item":"widget","amount":42.5}
{"order_id":"1002","item":"gadget","amount":19.99}' | kafkactl -C kafkactl.yaml --context vc produce orders
expected:
message: "2 messages produced"
return_code: 0
render_output: false
{% endvalidation %}
<!--vale on-->

## Validate

Consume the records straight from the broker to confirm that Kafka stores them as Avro, even though the producer sent JSON.
The `--print-schema` flag displays the Avro schema used for deserialization:

<!--vale off-->
{% validation custom-command %}
command: |
kafkactl -C kafkactl.yaml --context direct consume orders --from-beginning --exit --print-schema
expected:
message: '"order_id":"1001"'
return_code: 0
render_output: false
{% endvalidation %}
<!--vale on-->

Both records come back deserialized from Avro, alongside the schema used to decode them:

```shell
##{"type":"record","name":"Order","fields":[{"name":"order_id","type":"string"},{"name":"item","type":"string"},{"name":"amount","type":"double"}]}#1#{"order_id":"1001","item":"widget","amount":42.5}
##{"type":"record","name":"Order","fields":[{"name":"order_id","type":"string"},{"name":"item","type":"string"},{"name":"amount","type":"double"}]}#1#{"order_id":"1002","item":"gadget","amount":19.99}
```
{:.no-copy-code}

In this case, the producer never changed formats. It sent plain JSON, and the Schema Validation policy parsed it, then the Record Transcode Produce policy converted it to Avro before {{site.event_gateway_short}} wrote it to the backend cluster.
Loading