Webhook icon
Request icon
OutputValues icon
Log icon
If icon
TopicCreate icon
Produce icon
Consume icon
IonToJson icon
SlackIncomingWebhook icon
Fail icon

Block a breaking Avro schema before a Kafka producer deploys

Check a candidate Avro schema against Confluent Schema Registry, block breaking changes with field-level reasons, then register it and round trip a record.

Categories
DataInfrastructureinfrastructure

A breaking schema change is cheap to catch before a producer ships and expensive once it is on the topic. This flow puts the decision in front of the deploy. CI posts the candidate Avro schema to the webhook, the registry answers whether it is compatible with the subject's latest version. A FAILED execution is what stops the pipeline. When the answer is yes the flow registers the schema, writes one record with it, then reads that record back through the registry so the new version is proven usable rather than merely accepted.

The reason this is worth a flow rather than one curl is that the two registry endpoints it needs report failure in opposite ways:

  • POST /compatibility/subjects/{subject}/versions/latest answers HTTP 200 when the candidate is incompatible, with is_compatible: false in the body. A gate written against the status code passes every breaking change, silently. So gate branches on the parsed body field instead.
  • POST /subjects/{subject}/versions answers 409 Conflict on the same schema. There the status code is the right signal, so register_schema leaves allowFailed at its default false and lets the HTTP failure stop the run.

?verbose=true is load bearing, not cosmetic. The registry only fills in messages when it is set. Without it the flow can still block a deploy but cannot say which field broke, so the engineer gets a red pipeline with nothing to fix. With it the log carries the registry's own wording, for example The field 'currency' at path '/fields/1' in the new schema has no default value and is missing in the old schema.

The path uses the latest alias rather than a version number. That makes a brand new subject work with no special casing: the first candidate comes back compatible. A numeric version against a subject that does not exist is a version-not-found error instead of a verdict.

This blueprint was created by zkasuran.

How it works

  1. compatibility_level (io.kestra.plugin.core.http.Request) reads GET /config/{subject}?defaultToGlobal=true with allowFailed: true, so the log records whether the verdict came from the global level or a subject override. The registry applies the level itself, so this task is for the record rather than for the decision.
  2. check_compatibility posts the candidate to /compatibility/subjects/{subject}/versions/latest?verbose=true. The registry wants the schema as a JSON string inside the request body, which is what jq('{schemaType:"AVRO",schema:(.|tojson)}') builds. Pasting the schema in as a nested object is a 42201 that reads like a flow bug.
  3. verdict (io.kestra.plugin.core.output.OutputValues) names the three things the rest of the flow needs: the level, is_compatible plus the reasons joined into one line.
  4. gate (io.kestra.plugin.core.flow.If) branches on is_compatible from the body.
  5. Compatible, with register_on_pass on: register_schema registers the version, create_topic (io.kestra.plugin.kafka.TopicCreate) makes sure the topic exists, produce_record (io.kestra.plugin.kafka.Produce) writes one Avro record, consume_record (io.kestra.plugin.kafka.Consume) reads it back, readable_record converts the Ion output to JSON, then round_trip_report logs the new schema id with the record.
  6. Incompatible: log_blocked prints the registry reasons, alert_blocked posts them to Slack when notify_slack is on, then block_deploy (io.kestra.plugin.core.execution.Fail) ends the run as FAILED. Nothing is registered and nothing is produced.

Inputs

  • subject (STRING, default orders-value): the subject to check against.
  • candidate_schema (JSON): the schema about to be deployed.
  • sample_record (JSON): one record matching it, used for the round trip.
  • topic (STRING, default orders): topic for the round trip.
  • kafka_bootstrap_servers (STRING): broker list.
  • schema_registry_url (STRING): registry base URL with no trailing slash.
  • register_on_pass (BOOL, default true): set to false for a verdict-only dry run.
  • notify_slack (BOOL, default false): Slack alert on a block.

Prerequisites

  • A Kafka cluster plus a Confluent Schema Registry reachable from Kestra. A registry keeps its state in a Kafka topic, so it cannot run without a broker. This is why the blueprint is not marked as a demo.
  • Write permission on the subject for the registry user, plus produce and consume permission on the topic.
  • On a single broker test cluster, keep transactional: false on produce_record. The default transactional producer generates a transactional id that the default transaction.state.log.replication.factor of 3 cannot satisfy.
  • valueAvroSchema is supplied on produce_record alongside serdeProperties. The Avro serializer resolves its schema before any registry block is applied, so it needs the schema text as well as the registry URL.

Secrets

  • SCHEMA_GATE_WEBHOOK_KEY: the webhook key CI calls the flow with.
  • SLACK_WEBHOOK_URL: only needed when notify_slack is true.

Move kafka_bootstrap_servers plus schema_registry_url to secret() values if your registry sits behind credentials.

Quick start

  1. Add the flow, then set kafka_bootstrap_servers plus schema_registry_url to your own endpoints.
  2. Run it with the defaults against a subject that does not exist yet. The candidate comes back compatible, gets registered as version 1, then one record round trips.
  3. Run it again with a candidate that drops amount or adds a field with no default. The gate blocks, the log names the field, the execution ends FAILED.
  4. Point your CI at the webhook: curl -X POST -d '{"candidate_schema": ...}' <kestra>/api/v1/executions/webhook/company.team/kafka-schema-compatibility-gate/<key>. A non-zero result is the stop signal.

Expected outputs

  • outputs.verdict.values.is_compatible: true or false, read from the body rather than the status code.
  • outputs.verdict.values.reasons: the registry's errorType entries joined into one line, empty when compatible.
  • outputs.register_schema.body: {"id": <n>}, the global schema id of the newly registered version.
  • outputs.consume_record.messagesCount: 1 on a successful round trip.
  • A FAILED execution with the registry wording in the error message when the gate blocks.

How to extend

  • Swap schemaType to PROTOBUF or JSON to gate those contracts with the same two endpoints.
  • Check several subjects in one run by wrapping the gate in io.kestra.plugin.core.flow.Loop over a list of subject plus schema pairs.
  • Raise the level to FULL_TRANSITIVE for the subject with PUT /config/{subject} before the check, so consumers that lag behind are covered too.
  • Have the else branch open a pull request comment with the reasons instead of only failing.
  • Add io.kestra.plugin.core.trigger.Schedule beside the webhook to re-check the schema in git every night, which catches someone lowering the subject's compatibility level.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.