diff --git a/app/_event_gateway_policies/record-transcode-produce/index.md b/app/_event_gateway_policies/record-transcode-produce/index.md index 2bc2a0a1497..7ac9848b17d 100644 --- a/app/_event_gateway_policies/record-transcode-produce/index.md +++ b/app/_event_gateway_policies/record-transcode-produce/index.md @@ -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 %} diff --git a/app/_how-tos/event-gateway/consume-kafka-records-as-json-with-event-gateway.md b/app/_how-tos/event-gateway/consume-kafka-records-as-json-with-event-gateway.md index e3551a937c1..a857a888948 100644 --- a/app/_how-tos/event-gateway/consume-kafka-records-as-json-with-event-gateway.md +++ b/app/_how-tos/event-gateway/consume-kafka-records-as-json-with-event-gateway.md @@ -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 diff --git a/app/_how-tos/event-gateway/produce-kafka-records-as-avro-with-event-gateway.md b/app/_how-tos/event-gateway/produce-kafka-records-as-avro-with-event-gateway.md new file mode 100644 index 00000000000..1ac1bebd3dc --- /dev/null +++ b/app/_how-tos/event-gateway/produce-kafka-records-as-avro-with-event-gateway.md @@ -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
JSON] --> SV + + subgraph produce [Event Gateway Produce policy chain] + SV[Schema
Validation
Parse JSON] --> RT[Record Transcode
convert to Avro] + end + + RT --> K[Kafka
Broker
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: + + +{% 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 %} + + +## Create a listener with a forwarding policy + +Create a [listener](/event-gateway/entities/listener/) to accept connections: + + +{% 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 %} + + +Create a [Forward to Virtual Cluster policy](/event-gateway/policies/forward-to-virtual-cluster/) to forward traffic to the virtual cluster: + + +{% 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 %} + + +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`: + + +{% 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 %} + + +## 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`: + + +{% 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 %} + + +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 `-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: + + +{% 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 %} + + +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: + + +{% 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 %} + + +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: + + +{% validation custom-command %} +command: | + cat < 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 %} + + +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: + + +{% 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 %} + + +Produce two order records as plain JSON through the virtual cluster: + + +{% 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 %} + + +## 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: + + +{% 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 %} + + +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.