Skip to content

feat: support raw Avro values without Schema Registry - #338

Closed
YingqiDuan wants to merge 7 commits into
deviceinsight:mainfrom
YingqiDuan:feat/avro-schema-uri
Closed

YingqiDuan wants to merge 7 commits into
deviceinsight:mainfrom
YingqiDuan:feat/avro-schema-uri

Conversation

@YingqiDuan

Copy link
Copy Markdown

Description

kafkactl currently expects Avro messages to use Confluent framing, with each value prefixed by a magic byte and a Schema Registry ID. Some producers instead encode an Avro datum directly as the Kafka record value and provide the schema separately. CloudEvents messages in binary content mode are one example: the Kafka binding carries the schema reference in the ce_dataschema header.

This PR adds opt-in support for producing and consuming these raw Avro values. A standalone schema can be loaded from a local file or an HTTP(S) URL. When consuming, kafkactl can also read the schema URL from a Kafka header on each message.

No new dependencies or configuration keys are added.

Fixes #337

What changed

  • produce --avro-schema-file <path-or-url> encodes the record value as raw Avro binary, without Confluent framing.
  • consume --avro-schema-file <path-or-url> decodes record values with one fixed schema.
  • consume --avro-schema-header[=<name>] reads the schema URL from a Kafka message header. Without an explicit name, it uses ce_dataschema; pass --avro-schema-header=<name> to select a different header.
  • If both schema options are set, the fixed schema is loaded when the command starts, but the header takes priority when decoding a message. The fixed schema is selected only when the header is missing. If the header is present but invalid, the command returns an error instead of falling back.
  • Schemas resolved from message headers are cached by URL for the lifetime of the command. Set --avro-schema-cache=false to fetch and compile the schema again for each non-tombstone message that contains the header.
  • --print-schema works with standalone schemas, and the existing avro.jsonCodec setting continues to control whether standard JSON or Avro JSON is used.
  • Kafka tombstones are preserved, while valid zero-byte Avro values remain distinct from them.

Headers such as ce_dataschema and content-type can still be added with the existing -H/--header option. Raw Avro mode does not add them automatically.

Compatibility and scope

Raw Avro mode applies only to record values and is enabled only by the new raw Avro flags. When no raw Avro flag is set, codec selection follows the existing Schema Registry and default paths.

Raw Avro mode does not initialize the configured Schema Registry client. Schema HTTP requests are made only for URLs supplied through --avro-schema-file or the selected message header.

Plain keys continue to work, as do Protobuf keys backed by a local .proto file or protoset description. Registry-backed key codecs are not used while raw Avro mode is active.

Incompatible combinations are rejected before schema loading or Kafka client creation. These include raw Avro with --value-proto-type and, on produce, non-default --key-schema-version or --value-schema-version settings.

This PR does not inspect content-type, detect raw Avro automatically, add CloudEvents headers, perform full CloudEvents validation, or support CloudEvents structured content mode.

Remote schema safety

Header-based lookup is opt-in because the target URL comes from Kafka message data. Remote schema loading:

  • accepts only HTTP(S) URLs and requires a final 2xx response;
  • uses a five-second HTTP client timeout and limits the response body to 1 MiB;
  • follows at most ten redirects and validates every redirect target;
  • rejects URLs containing user information;
  • does not add Schema Registry credentials, cookies, or authentication headers to schema requests; and
  • keeps URL user information, query strings, response bodies, and untrusted redirect targets out of error messages.

Localhost and private-network addresses are allowed so local and in-cluster schema services can be used. A message header can therefore trigger requests to services reachable by kafkactl. Enable header-based lookup only when every producer that can write to the topic is trusted.

The header URL cache has no entry limit. It can be disabled for topics that use mutable URLs or many distinct schema URLs.

Related hardening

  • Caps the item count for a single Avro array or map block at 1,000,000 during decoding. This process-wide limit also applies to Registry-backed Avro decoding, but it does not limit the cumulative number of items across multiple blocks.
  • Decodes AMQP str32 and vbin32 lengths safely on 32-bit systems and handles valid zero-length values correctly. The fix also applies to schema URL headers received through Kafka-compatible services such as Azure Event Hubs.

Testing

The relevant package tests passed without cached results:

go test -count=1 -short \
  ./internal/helpers/avro \
  ./internal/consume \
  ./internal/producer \
  ./cmd/consume \
  ./cmd/produce

go test -race also passed for the same packages in a Linux Docker environment. go vet, go tool golangci-lint run, go build ./..., and the AMQP tests with GOARCH=386 also passed.

The following integration tests passed against ZooKeeper and three Kafka brokers without a running Schema Registry:

  • TestConsumeRawAvroIntegration
  • TestConsumeRawAvroProtobufKeyIntegration
  • TestProduceRawAvroNullIntegration
  • TestProduceRawAvroProtobufKeyIntegration

These tests cover per-message schemas, fixed-schema fallback, caching on and off, tombstones, zero-byte Avro values, and local Protobuf keys. They also verify that raw Avro values can be produced and consumed without Schema Registry.

Type of change

  • Bug fix (non-breaking change which fixes an issue)
  • New feature (non-breaking change which adds functionality)

Documentation

  • the change is mentioned in the ## [Unreleased] section of CHANGELOG.md
  • a usage example was added to README.adoc
  • tests for the changes have been implemented (see: Testing your changes)

@YingqiDuan YingqiDuan closed this Sep 6, 2026
@YingqiDuan
YingqiDuan deleted the feat/avro-schema-uri branch September 6, 2026 18:26
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support Avro decode/encode via a per-message schema-URI header (CloudEvents dataschema), not just Confluent Schema Registry

1 participant